Spring Integration问题:MessageSource不遵循errorChannel头信息
S3StreamingMessageSource Respect the errorChannel Header Let's break down the problem first: Your current flow enriches messages with the S3_ERROR_CHANNEL header, but when you invoke operations on S3StreamingMessageSource inside the handle method, any exceptions thrown don't route to this custom error channel. That's because the messageSource's operations run synchronously within the handler, and default exception handling won't automatically pick up the header unless you explicitly wire it in.
Here are concrete, code-aligned solutions to fix this:
Solution 1: Explicitly Catch and Route Exceptions to the Custom Error Channel
Wrap your S3StreamingMessageSource operations in a try-catch block, then manually send any caught exceptions to your S3_ERROR_CHANNEL using a injected MessageChannel bean.
First, add the error channel injection to your component:
@Resource(name = S3_ERROR_CHANNEL) private MessageChannel s3ErrorChannel; @Resource(name = S3_CLIENT_BEAN) private MessageSource<InputStream> messageSource;
Then update your handle logic to handle exceptions explicitly:
public IntegrationFlow fileStreamingFlow() { return IntegrationFlows.from(s3Properties.getFileStreamingInputChannel()) .enrichHeaders(spec -> spec.header(ERROR_CHANNEL, S3_ERROR_CHANNEL, true)) .handle(String.class, (fileName, headers) -> { if (messageSource instanceof S3StreamingMessageSource) { S3StreamingMessageSource s3StreamingMessageSource = (S3StreamingMessageSource) messageSource; try { // Your existing logic with s3StreamingMessageSource goes here Message<InputStream> s3Message = s3StreamingMessageSource.receive(); // Process the input stream... } catch (Exception e) { // Build an error message with original headers and the exception ErrorMessage errorMessage = new ErrorMessage(e, headers); s3ErrorChannel.send(errorMessage); // Optional: rethrow to trigger downstream error handling or return an error marker throw new RuntimeException("Failed to process S3 file: " + fileName, e); } } // Rest of your handling logic return null; }) .get(); }
Solution 2: Use Spring Integration's RequestHandlerRetryAdvice for Automated Error Routing
Add a retry/error handling advice to your handler, which will automatically route exceptions to the errorChannel header value. This is useful if you want to add retry logic alongside error routing.
First, configure the advice bean:
@Bean public RequestHandlerRetryAdvice s3RetryAdvice() { RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice(); // Optional: Configure retry settings (adjust max attempts as needed) RetryTemplate retryTemplate = new RetryTemplate(); SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); retryPolicy.setMaxAttempts(3); retryTemplate.setRetryPolicy(retryPolicy); advice.setRetryTemplate(retryTemplate); // Set a recoverer that uses the error channel from the message headers advice.setRecoveryCallback(new ErrorMessageSendingRecoverer(s3ErrorChannel) { @Override protected Object createErrorMessage(Throwable t, Message<?> message) { // Prioritize the error channel from the original message headers MessageChannel targetChannel = (MessageChannel) message.getHeaders().get(ERROR_CHANNEL); if (targetChannel != null) { return new ErrorMessage(t, message.getHeaders()); } return super.createErrorMessage(t, message); } }); return advice; }
Then attach the advice to your handler in the flow:
public IntegrationFlow fileStreamingFlow() { return IntegrationFlows.from(s3Properties.getFileStreamingInputChannel()) .enrichHeaders(spec -> spec.header(ERROR_CHANNEL, S3_ERROR_CHANNEL, true)) .handle(String.class, (fileName, headers) -> { if (messageSource instanceof S3StreamingMessageSource) { S3StreamingMessageSource s3StreamingMessageSource = (S3StreamingMessageSource) messageSource; // Your existing logic here } // Rest of your handling logic return null; }, handlerSpec -> handlerSpec.advice(s3RetryAdvice())) // Attach the retry/error advice .get(); }
Solution 3: Propagate Message Context to S3StreamingMessageSource (Spring Integration 5.3+)
If you're using Spring Integration 5.3 or later, use MessageRequestHandlerAdvice to propagate the current message context to the S3StreamingMessageSource operations. This ensures exceptions respect the errorChannel header without manual exception handling.
First, define the advice bean:
@Bean public MessageRequestHandlerAdvice messageContextAdvice() { return new MessageRequestHandlerAdvice(); }
Then update your flow to use the advice:
public IntegrationFlow fileStreamingFlow() { return IntegrationFlows.from(s3Properties.getFileStreamingInputChannel()) .enrichHeaders(spec -> spec.header(ERROR_CHANNEL, S3_ERROR_CHANNEL, true)) .handle(String.class, (fileName, headers) -> { if (messageSource instanceof S3StreamingMessageSource) { S3StreamingMessageSource s3StreamingMessageSource = (S3StreamingMessageSource) messageSource; // Your existing logic here } // Rest of your handling logic return null; }, handlerSpec -> handlerSpec.advice(messageContextAdvice())) .get(); }
Key Context
The core issue is that S3StreamingMessageSource operations don't automatically inherit the current message's headers unless you explicitly propagate the context or handle exceptions manually. If you were using S3StreamingMessageSource as the direct source of the flow (e.g., IntegrationFlows.from(messageSource, configurer -> configurer.poller(...))), you'd configure the poller's error channel directly via pollerSpec.errorChannel(S3_ERROR_CHANNEL). But since you're invoking it within a handler, the solutions above are the right fit.
内容的提问来源于stack exchange,提问作者Abhijit Sarkar

