You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.19 04:19:52