求基于Java注解/DSL的聚合器超时释放策略示例
刚好我之前做过类似的需求,针对你要的以超时作为释放策略的非XML(Java注解/DSL)实现,这里给你几个基于Spring Integration的实用示例(毕竟你提到的TimeoutCountSequenceSizeReleaseStrategy就是Spring Integration里的核心类):
基于Java注解的完整配置示例
如果你用Spring Integration,可以通过纯注解配置实现“最后一个关联条目添加后延迟一段时间释放”的逻辑:
- 首先创建配置类,实例化超时释放策略、消息存储和聚合器:
@Configuration @EnableIntegration public class AggregatorConfig { // 配置超时释放策略:最后一条消息到达后延迟5秒触发释放 // 参数依次是:超时时间(毫秒)、计数阈值、序列大小阈值 // 这里我们只关注超时逻辑,所以把计数和序列阈值设为Integer.MAX_VALUE,只触发超时规则 @Bean public ReleaseStrategy timeoutReleaseStrategy() { return new TimeoutCountSequenceSizeReleaseStrategy(5000, Integer.MAX_VALUE, Integer.MAX_VALUE); } // 配置支持超时的消息组存储,这里用SimpleMessageStore(生产环境可以用RedisMessageStore做分布式场景) @Bean public MessageGroupStore messageGroupStore() { SimpleMessageStore messageStore = new SimpleMessageStore(); // 设置消息组的超时时间,和释放策略的超时保持一致 messageStore.setTimeout(5000); return messageStore; } // 定义聚合器处理器 @Bean public MessageHandler aggregatorHandler() { AggregatingMessageHandler aggregator = new AggregatingMessageHandler( new DefaultAggregatingMessageGroupProcessor(), messageGroupStore() ); aggregator.setReleaseStrategy(timeoutReleaseStrategy()); // 开启超时后发送已聚合的部分结果(如果需要) aggregator.setSendPartialResultOnExpiry(true); return aggregator; } // 配置输入输出通道 @Bean public DirectChannel inputChannel() { return new DirectChannel(); } @Bean public DirectChannel outputChannel() { return new DirectChannel(); } // 绑定聚合器到通道流程 @Bean public IntegrationFlow aggregatorFlow() { return IntegrationFlows.from(inputChannel()) .handle(aggregatorHandler()) .channel(outputChannel()) .get(); } }
- 业务代码中发送消息,注意给同组消息设置相同的
correlationId:
@Autowired private MessageChannel inputChannel; public void sendGroupedMessages() { // 发送同组的多条消息 inputChannel.send(MessageBuilder.withPayload("关联条目1") .setCorrelationId("group-001") .build()); inputChannel.send(MessageBuilder.withPayload("关联条目2") .setCorrelationId("group-001") .build()); // 最后一条消息发送后,等待5秒就会触发聚合释放 inputChannel.send(MessageBuilder.withPayload("关联条目3(最后一条)") .setCorrelationId("group-001") .build()); }
基于Spring Integration DSL的简化实现
如果想用更简洁的DSL风格配置,代码会更紧凑:
@Configuration @EnableIntegration public class AggregatorDslConfig { @Bean public IntegrationFlow timeoutAggregatorFlow() { return IntegrationFlows.from("inputChannel") .aggregate(aggregatorSpec -> aggregatorSpec // 指定超时释放策略 .releaseStrategy(new TimeoutCountSequenceSizeReleaseStrategy(5000, Integer.MAX_VALUE, Integer.MAX_VALUE)) // 绑定消息存储 .messageStore(messageGroupStore()) // 超时后发送部分结果 .sendPartialResultOnExpiry(true) // 按correlationId分组 .correlationStrategy(msg -> msg.getHeaders().get("correlationId"))) .channel("outputChannel") .get(); } @Bean public MessageGroupStore messageGroupStore() { SimpleMessageStore store = new SimpleMessageStore(); store.setTimeout(5000); return store; } }
自定义超时释放策略(更灵活的场景)
如果默认的TimeoutCountSequenceSizeReleaseStrategy不能满足你的特殊逻辑,比如需要判断“是否是最后一条消息”再触发超时,可以自定义ReleaseStrategy:
public class CustomTimeoutReleaseStrategy implements ReleaseStrategy { private final long timeoutMillis; public CustomTimeoutReleaseStrategy(long timeoutMillis) { this.timeoutMillis = timeoutMillis; } @Override public boolean canRelease(MessageGroup messageGroup) { // 先判断该消息组是否被标记为"已完成"(发送最后一条消息时设置标记) boolean isGroupCompleted = messageGroup.isComplete(); // 计算最后一条消息到现在的时间差 long timeElapsedSinceLastMsg = System.currentTimeMillis() - messageGroup.getLastModified(); // 满足"已完成"且超时条件时,允许释放 return isGroupCompleted && timeElapsedSinceLastMsg >= timeoutMillis; } }
使用时,在聚合器配置里替换成这个自定义策略,并且发送最后一条消息时标记消息组完成:
// 发送最后一条消息时,设置COMPLETE_GROUP_HEADER为true inputChannel.send(MessageBuilder.withPayload("最后一条关联条目") .setCorrelationId("group-001") .setHeader(MessageHeaders.COMPLETE_GROUP_HEADER, true) .build());
内容的提问来源于stack exchange,提问作者max3d
相关产品推荐
相关产品推荐

