Spring Integration AWS Kinesis聚合器仅首次触发后续消息不处理问题排查
问题根因
Spring Integration 聚合器默认行为导致的问题,核心有两点:
- 聚合器默认将触发过释放策略的消息组标记为
COMPLETED状态,后续携带相同correlation key(即你配置的headers['foo']同值)的消息会被直接丢弃,不会进入聚合流程。你当前设置的释放策略阈值为1,第一条消息到达就会触发组释放,同key的后续消息自然无法触发后续流程。 - 你未配置专用MessageStore,默认使用的
SimpleMessageStore不会自动清理已完成的消息组,已标记完成的组会一直留存,导致后续同key消息持续被拦截。
解决方案
只需在聚合器配置中新增一行,开启「完成后自动移除消息组」的配置即可:
@Bean @ServiceActivator(inputChannel = "kinesisSendChannel") public MessageHandler aggregator(TestMessageProcessor messageProcessor, MessageChannel resultChannel, TimeoutCountSequenceSizeReleaseStrategy timeoutCountSequenceSizeReleaseStrategy) { AggregatingMessageHandler handler = new AggregatingMessageHandler(messageProcessor); handler.setCorrelationStrategy(new ExpressionEvaluatingCorrelationStrategy("headers['foo']")); handler.setReleaseStrategy(timeoutCountSequenceSizeReleaseStrategy); handler.setOutputProcessor(messageProcessor); handler.setOutputChannel(resultChannel); // 新增以下配置,组完成后直接从存储中移除,后续同key消息会新建组处理 handler.setExpireGroupsUponCompletion(true); return handler; }
如果你的业务需要按超时分批聚合同key消息,可以额外添加handler.setSendPartialResultOnExpiry(true)配置,配合你设置的60秒超时策略使用,超时后自动释放当前组内的消息并清理组。
内容的提问来源于stack exchange,提问作者LBPS
相关产品推荐
相关产品推荐

