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

Spring Integration应用重启后如何发送过期待处理消息组?

Spring Integration 重启后过期消息组无法发送的问题与解决

我开发了一个基于Spring Integration的小型应用,消息和消息组存储在数据库中。现在遇到一个问题:部分消息/消息组在等待超时触发发送时,应用意外重启。重启后数据库里还保留着这些消息,但它们无法被自动发送出去。我需要配置来处理这些数据库中的过期消息组,或者恢复对应的超时定时器。尝试过使用MessageGroupStoreReaper,但没达到预期效果。

原配置代码

@Configuration
public class ConsumingChannelConfig {

    @Bean
    public DirectChannel consumingChannel() {
        return new DirectChannel();
    }

    @Bean
    public KafkaMessageDrivenChannelAdapter<String, String> kafkaMessageDrivenChannelAdapter() {
        KafkaMessageDrivenChannelAdapter<String, String> kafkaMessageDrivenChannelAdapter =
                new KafkaMessageDrivenChannelAdapter<>(kafkaListenerContainer());
        kafkaMessageDrivenChannelAdapter.setOutputChannel(consumingChannel());
        MessagingMessageConverter messageConverter = new MessagingMessageConverter();
        messageConverter.setGenerateMessageId(true);
        kafkaMessageDrivenChannelAdapter.setRecordMessageConverter(messageConverter);
        return kafkaMessageDrivenChannelAdapter;
    }

    @Bean
    public DataSource getDataSource() {
        return ...;
    }

    @Bean
    public JdbcMessageStore jdbcMessageStore() {
        return new JdbcMessageStore(getDataSource());
    }

    @ServiceActivator(inputChannel = "consumingChannel")
    @Bean
    public MessageHandler aggregator() {
        long timeout = 10000L;
        AggregatingMessageHandler aggregator =
                new AggregatingMessageHandler(new DefaultAggregatingMessageGroupProcessor(),
                        jdbcMessageStore());
        aggregator.setOutputChannel((message, l) -> {
            System.out.println("MESSAGE: " + message);
            return true;
        });
        aggregator.setGroupTimeoutExpression(new ValueExpression<>(timeout));
//        aggregator.setTaskScheduler(this.taskScheduler);
        aggregator.setCorrelationStrategy(new MyCorrelationStrategy());
        aggregator.setSendPartialResultOnExpiry(true);
        aggregator.setExpireGroupsUponCompletion(true);
        aggregator.setExpireGroupsUponTimeout(true);
        aggregator.setDiscardChannel((message, timeout1) -> {
            System.out.println("DISCARD: " + message + ", timeout: " + timeout1);
            return true;
        });
        aggregator.setReleaseStrategy(new ReleaseStrategy() {
            @Override
            public boolean canRelease(MessageGroup group) {
                return System.currentTimeMillis() - group.getTimestamp() >= timeout;
            }
        });
        return aggregator;
    }

    @Bean
    public MessageGroupStoreReaper reaper() {
        MessageGroupStoreReaper reaper = new MessageGroupStoreReaper(jdbcMessageStore());
        reaper.setPhase(1);
        reaper.setTimeout(2000L);
        reaper.setAutoStartup(true);
//        reaper.setExpireOnDestroy(true);
        return reaper;
    }

    @Bean
    public ConcurrentMessageListenerContainer<String, String> kafkaListenerContainer() {
        ContainerProperties containerProps = new ContainerProperties("spring-integration-topic");

        return new ConcurrentMessageListenerContainer<>(
                consumerFactory(), containerProps);
    }

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> properties = new HashMap<>();
        properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
        properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        properties.put(ConsumerConfig.GROUP_ID_CONFIG, "spring-integration");
        // automatically reset the offset to the earliest offset
        properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
//        DefaultKafkaHeaderMapper mapper = new DefaultKafkaHeaderMapper();
        return properties;
    }
}

解决方案

通过启用Spring调度,在应用启动后手动触发一次MessageGroupStoreReaper的执行,处理数据库中遗留的过期消息组:

@EnableScheduling
@SpringBootApplication
public class SpringIntegrationExampleApplication {

    public static void main(String[] args) {
        SpringApplication.run(SpringIntegrationExampleApplication.class, args);
    }


    @Autowired
    private MessageGroupStoreReaper reaper;

    @Scheduled(initialDelay = 2000, fixedDelay = Long.MAX_VALUE)
    public void start() {
        reaper.run();
    }

}

内容的提问来源于stack exchange,提问作者Oleksandr H

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 02:17:04