Spring Integration处理IBM MQ批量地址消息的异常处理优化诉求
问题场景与现有代码
背景
我有一条来自IBM MQ的消息,消息内包含地址列表,通过Spring Integration Flow处理,但遇到以下问题:DAO处理单个地址时抛出异常(如SQLException)会传递至父通道,导致消息被提交回MQ并触发重试。期望完成列表中所有地址的处理后,再将错误消息提交回MQ。此前尝试使用错误通道和容器级错误处理器,但未生效。
现有代码实现
ListenerContainer
public DefaultMessageListenerContainer addressProcessorListenerContainer( ConnectionFactory connectionFactory, String queueName, int consumers, PipelineErrorHandler pipelineErrorHandler ) { DefaultMessageListenerContainer container = new DefaultMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setSessionTransacted(true); container.setPubSubDomain(false); container.setConcurrentConsumers(consumers); container.setDestinationName(queueName); container.setErrorHandler(pipelineErrorHandler); return container; }
Error Handler
@Component public class PipelineErrorHandler implements ErrorHandler { private static final Logger LOGGER = LoggerFactory.getLogger(PipelineErrorHandler.class); @Override public void handleError(Throwable t) { LOGGER.error("Error occurred during pipeline execution, It will retry now", t); } }
父IntegrationFlow
@Bean public IntegrationFlow processAddressFlow(DefaultMessageListenerContainer addressProcessorListenerContainer) { return IntegrationFlows.from(Jms.messageDrivenChannelAdapter(addressProcessorListenerContainer)) //Router to decide which channel to follow based on function .route(Message.class, this::hasMessageBreachedItsMaxRedeliveryThreshold, router -> router.channelMapping(true, "max-redelivery-threshold-channel").applySequence(true) .channelMapping(false, "process-address-channel")) .get(); }
处理地址的通道
@Bean public IntegrationFlow processAddressEventFlow( EventTypeFilter eventTypeFilter, ) { return IntegrationFlows.from("process-address-events-channel") .filter(eventTypeFilter, spec -> spec.id("processAddressEventFlow.eventTypeFilter")) .split(spec -> spec.id("processAddressEventFlow.split")) .channel("process-address-event-channel") .get(); }
该通道会拆分地址列表并传递至process-address-event-channel通道。
调用DAO更新地址至数据库的最终通道
public IntegrationFlow addressUpdateFlow(AddressUpdateDao addressUpdateDao, MessageSourceSystemFilter messageSourceSystemFilter) { return IntegrationFlows.from("process-address-event-channel") .handle(AddressUpdateDao.class, addressUpdateDao, spec -> spec.id("addressUpdateFlow.addressUpdateFlow")) .log() .get(); }
解决方案
要实现“处理完所有地址后再决定是否重试消息”的目标,核心是将拆分后的子消息的错误集中处理,待所有子消息处理完成后再判断是否触发父消息的重试,可以通过以下步骤实现:
1. 为拆分后的子消息配置局部错误处理
在addressUpdateFlow的handle步骤添加错误通道,捕获单个地址处理的异常,避免异常直接冒泡到父容器:
public IntegrationFlow addressUpdateFlow(AddressUpdateDao addressUpdateDao, MessageSourceSystemFilter messageSourceSystemFilter) { return IntegrationFlows.from("process-address-event-channel") .handle(AddressUpdateDao.class, addressUpdateDao, spec -> spec .id("addressUpdateFlow.addressUpdateFlow") .errorChannel("address-processing-error-channel")) // 指定局部错误通道 .log() .get(); }
2. 实现错误收集组件,记录子消息错误
创建一个组件用于跟踪同一父消息下的子处理错误:
@Component public class AddressErrorCollector { private static final Logger LOGGER = LoggerFactory.getLogger(AddressErrorCollector.class); private final Map<String, AtomicInteger> errorCountByCorrelationId = new ConcurrentHashMap<>(); public void collectError(ErrorMessage errorMessage) { String correlationId = errorMessage.getHeaders().get(IntegrationMessageHeaderAccessor.CORRELATION_ID, String.class); if (correlationId != null) { errorCountByCorrelationId.computeIfAbsent(correlationId, k -> new AtomicInteger(0)).incrementAndGet(); } LOGGER.error("Failed to process address, correlationId: {}", correlationId, errorMessage.getPayload()); } public int getErrorCount(String correlationId) { return errorCountByCorrelationId.getOrDefault(correlationId, new AtomicInteger(0)).get(); } public void cleanUp(String correlationId) { errorCountByCorrelationId.remove(correlationId); } }
配置错误通道的处理Flow:
@Bean public IntegrationFlow addressProcessingErrorFlow(AddressErrorCollector errorCollector) { return IntegrationFlows.from("address-processing-error-channel") .handle(errorCollector, "collectError") .get(); }
3. 使用聚合器等待所有子消息处理完成,统一判断重试
修改processAddressEventFlow,在拆分后添加聚合器,等待所有子消息处理完毕后,根据错误计数决定是否抛出异常触发MQ重试:
@Bean public IntegrationFlow processAddressEventFlow(EventTypeFilter eventTypeFilter, AddressErrorCollector errorCollector) { return IntegrationFlows.from("process-address-events-channel") .filter(eventTypeFilter, spec -> spec.id("processAddressEventFlow.eventTypeFilter")) .split(spec -> spec.id("processAddressEventFlow.split").applySequence(true)) // 开启序列标记,让聚合器识别子消息归属 .channel("process-address-event-channel") .aggregate(spec -> spec .id("processAddressEventFlow.aggregator") .correlationStrategy(m -> m.getHeaders().get(IntegrationMessageHeaderAccessor.CORRELATION_ID)) .releaseStrategy(group -> group.isComplete()) // 所有子消息处理完成后释放 .outputProcessor(group -> { String correlationId = group.getCorrelationKey().toString(); int errorCount = errorCollector.getErrorCount(correlationId); errorCollector.cleanUp(correlationId); // 清理缓存,避免内存泄漏 if (errorCount > 0) { // 抛出异常,触发JMS容器的重试逻辑 throw new RuntimeException(String.format("Failed to process %d addresses in message", errorCount)); } return group.getMessages(); })) .get(); }
4. 容器错误处理器保持原有逻辑
原PipelineErrorHandler无需修改,只有当聚合器抛出异常时,才会触发容器的重试逻辑,此时整个消息的所有地址都已处理完毕。
关键原理说明
- 拆分+聚合:通过
applySequence(true)为子消息添加关联ID和序列信息,聚合器可据此识别同一父消息的所有子消息,等待全部处理完成后执行统一判断。 - 局部错误通道:拦截单个子消息的异常,避免直接触发父消息重试,先收集错误信息。
- 错误计数与清理:通过
AddressErrorCollector跟踪每个父消息的错误数量,聚合完成后清理缓存,防止内存泄漏。
内容的提问来源于stack exchange,提问作者Prateek Sharma
相关产品推荐
相关产品推荐

