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
相关产品推荐
相关产品推荐

