【问题标题】:Spring Integration: MessageSource doesn't honor errorChannel headerSpring Integration:MessageSource 不支持 errorChannel 标头
【发布时间】:2018-04-19 00:56:15
【问题描述】:

我有以下流程:

@Resource(name = S3_CLIENT_BEAN)
private MessageSource<InputStream> messageSource;

public IntegrationFlow fileStreamingFlow() {
    return IntegrationFlows.from(s3Properties.getFileStreamingInputChannel())
        .enrichHeaders(spec -> spec.header(ERROR_CHANNEL, S3_ERROR_CHANNEL, true))
        .handle(String.class, (fileName, h) -> {
                    if (messageSource instanceof S3StreamingMessageSource) {
                        S3StreamingMessageSource s3StreamingMessageSource = (S3StreamingMessageSource) messageSource;

                        ChainFileListFilter<S3ObjectSummary> chainFileListFilter = new ChainFileListFilter<>();
                        chainFileListFilter.addFilters(...);
                        s3StreamingMessageSource.setFilter(chainFileListFilter);

                        return s3StreamingMessageSource.receive();
                    }                    
                    return messageSource.receive();
                }, spec -> spec
                        .requiresReply(false) // in case all messages got filtered out
        )
        .channel(s3Properties.getFileStreamingOutputChannel())
        .get();
}

我发现如果s3StreamingMessageSource.receive 抛出异常,错误最终会出现在为管道中的前一个流配置的错误通道中,而不是为此流配置的S3_ERROR_CHANNEL。不确定是否与this 问题有关。

【问题讨论】:

    标签: spring-integration spring-integration-aws


    【解决方案1】:

    s3StreamingMessageSource.receive() 是从SourcePollingChannelAdapter 调用的:

    protected Message<?> receiveMessage() {
        return this.source.receive();
    }
    

    这个来自AbstractPollingEndpoint:

    private boolean doPoll() {
    
            message = this.receiveMessage();
    ...
    
            this.handleMessage(message);
    ...
    }
    

    handleMessage() 这样做:

    this.messagingTemplate.send(getOutputChannel(), message);
    

    所以,那肯定离下游提到的.enrichHeaders(spec -&gt; spec.header(ERROR_CHANNEL, S3_ERROR_CHANNEL, true))还很远。

    但是,您仍然可以在该 S3_ERROR_CHANNEL 中捕获异常。注意IntegrationFlows.from()的第二个参数:

    IntegrationFlows.from(s3Properties.getFileStreamingInputChannel(),
               e -> e.poller(Pollers.fixedDelay(...)
                           .errorChannel(S3_ERROR_CHANNEL)))
    

    或者,根据您当前的情况,您在某个地方有一个全局轮询器,所以在那里配置一个errorChannel。

    【讨论】:

    • 我没有投票器。我可以使用建议或MessagePublishingErrorHandler 做某事吗?
    • 嗯,在你的应用程序上下文中肯定有一个轮询器。否则MessageSource 不起作用。如果你不能修改那个,只需使用我建议的 API 配置一个本地的。是的,您可以考虑与上述发布者一起实施一些建议,但如果轮询器上的 errorChannel 已经为您完成了逻辑,为什么还要重新发明轮子
    • 某处可能有一个轮询器(不在代码附近),但与显示的流程无关。显示的流程是通过将消息放入输入通道来触发的。见this
    • 好的!您必须在问题中提到您手动调用MessageSource.receive()。因此,receive() 的异常仍然远离提到的标头丰富器,这是很自然的。所以,是的,你必须手动调用try...catch 并确保手动发送错误消息。
    • 您对如何手动发送错误消息有任何建议/示例吗?
    猜你喜欢
    • 1970-01-01
    • 2023-03-20
    • 1970-01-01
    • 1970-01-01
    • 2017-11-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-13
    相关资源
    最近更新 更多