Spring Integration中Poller与Aggregator结合使用遇到问题求助
Spring Integration Aggregator 未激活问题解决指南
问题背景
我在使用Spring Integration Aggregator时遇到了问题:当前通过Poller从SMB文件系统拉取文件并解析XML的流程正常,但想要将单次轮询得到的所有解析后XML文件聚合为一个单元处理,虽然通过Poller Advice给每个消息硬编码了CORRELATION_ID,但Aggregator始终未被触发。
代码示例
@Configuration public class SMBConfig { /** * SMB Session Factory maintains persistent connection to SMB Share */ @Bean public SessionFactory<jcifs.smb.SmbFile> cachingSessionFactory() { CIFSContext context = SingletonContext.getInstance().withAnonymousCredentials(); SmbSessionFactory smbSessionFactory = new SmbSessionFactory(context); smbSessionFactory.setHost(host); smbSessionFactory.setPort(port); smbSessionFactory.setDomain(domain); smbSessionFactory.setShareAndDir(share); smbSessionFactory.setSmbMinVersion(DialectVersion.SMB210); smbSessionFactory.setSmbMaxVersion(DialectVersion.SMB311); return new CachingSessionFactory<>(smbSessionFactory); } @Bean public SmbRemoteFileTemplate template(SessionFactory<jcifs.smb.SmbFile> cachingSessionFactory) { return new SmbRemoteFileTemplate(cachingSessionFactory); } /** * Transforms a Stream into a parsed and typed Event */ @Bean @Transformer(inputChannel = "stream", outputChannel = "xml") public org.springframework.integration.transformer.Transformer transformer(XmlMapper xmlMapper) { return new XmlStreamTransformer("UTF-8", xmlMapper); } @Bean @InboundChannelAdapter(channel = "stream", poller = @Poller(value = "pollerMetadata")) public MessageSource<InputStream> smbMessageSource(CompositeFileListFilter<SmbFile> filter, SmbRemoteFileTemplate template) { SmbStreamingMessageSource messageSource = new SmbStreamingMessageSource(template); messageSource.setRemoteDirectory(dir); messageSource.setMaxFetchSize(2); return messageSource; } @Bean public PollerMetadata pollerMetadata() { return Pollers.fixedDelay(30000) .maxMessagesPerPoll(10000) .advice(new ReceiveMessageAdvice() { @Override public Message<?> afterReceive(Message<?> result, Object source) { return result == null ? null : MessageBuilder.fromMessage(result) .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, "123") .build(); } }) .getObject(); } } @MessageEndpoint @Slf4j class EventAggregator { @Aggregator(inputChannel = "xml") public Long aggregatingMethod(List<EventWrapper> events) { log.info("SAggregator"); return 1L; } }
问题根源与解决方法
1. 核心问题:缺少聚合触发条件
Aggregator仅收到CORRELATION_ID时,无法判断何时完成聚合。默认逻辑需要消息携带SEQUENCE_NUMBER(消息序号)和SEQUENCE_SIZE(批次总数量),或者通过超时、自定义策略触发聚合。
2. 针对单次轮询的聚合方案
方案一:设置聚合组超时(简单高效)
直接为Aggregator配置groupTimeout,当指定时间内没有同组新消息到达时,自动触发聚合。修改EventAggregator的注解即可:
@MessageEndpoint @Slf4j class EventAggregator { // 添加groupTimeout,轮询结束后1秒触发聚合(时间可根据实际调整) @Aggregator(inputChannel = "xml", groupTimeout = "1000") public Long aggregatingMethod(List<EventWrapper> events) { log.info("Aggregator triggered, processed {} events", events.size()); return (long) events.size(); } }
同时建议为每次轮询生成唯一的CORRELATION_ID(比如用时间戳),避免不同轮询的消息被错误聚合:
@Bean public PollerMetadata pollerMetadata() { return Pollers.fixedDelay(30000) .maxMessagesPerPoll(10000) .advice(new ReceiveMessageAdvice() { @Override public Message<?> afterReceive(Message<?> result, Object source) { if (result == null) return null; // 用当前时间戳生成唯一批次ID String batchId = "poll-batch-" + System.currentTimeMillis(); return MessageBuilder.fromMessage(result) .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, batchId) .build(); } }) .getObject(); }
方案二:基于序号触发聚合(精确控制)
如果需要严格在单次轮询的所有消息到达后立即聚合,可以在Poller Advice中为消息添加SEQUENCE_NUMBER和SEQUENCE_SIZE:
@Bean public PollerMetadata pollerMetadata() { AtomicInteger sequenceCounter = new AtomicInteger(0); AtomicInteger batchSize = new AtomicInteger(0); return Pollers.fixedDelay(30000) .maxMessagesPerPoll(10000) .advice(new ReceiveMessageAdvice() { // 轮询开始前获取本次要拉取的文件总数 @Override public Message<?> beforeReceive(Object source) { sequenceCounter.set(0); if (source instanceof SmbStreamingMessageSource) { try { List<SmbFile> files = ((SmbStreamingMessageSource) source).listFiles(); batchSize.set(files.size()); } catch (Exception e) { batchSize.set(0); } } return null; } // 为每个消息设置批次ID、序号和总数量 @Override public Message<?> afterReceive(Message<?> result, Object source) { if (result == null) return null; int seqNum = sequenceCounter.incrementAndGet(); String batchId = "poll-batch-" + System.currentTimeMillis(); return MessageBuilder.fromMessage(result) .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, batchId) .setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, seqNum) .setHeader(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE, batchSize.get()) .build(); } }) .getObject(); }
此时当最后一条消息(SEQUENCE_NUMBER等于SEQUENCE_SIZE)到达Aggregator时,会立即触发聚合。
3. 额外检查点
- 确保
EventAggregator类被Spring正确扫描(@MessageEndpoint生效,或在配置类中注册为Bean)。 - 确认
xml通道配置正确,解析后的消息能正常发送到该通道。
内容的提问来源于stack exchange,提问作者jdh961502
相关产品推荐
相关产品推荐

