Spring Integration DSL中配置幂等接收器处理JMS重复消息
问题解答
1. 在JMS消息驱动通道适配器中配置幂等拦截器
你可以直接在IntegrationFlow的处理链中通过interceptor()方法添加自定义的IdempotentReceiverInterceptor——因为它实现了ChannelInterceptor接口,能直接作用于消息流。修改后的代码如下:
.from(Jms.messageDrivenChannelAdapter(activeMQConnectionFactory) .destination(sourceQueue) .configureListenerContainer(spec -> { spec.sessionTransacted(false); spec.sessionAcknowledgeMode(Session.DUPS_OK_ACKNOWLEDGE); })) .interceptor(idempotentReceiverInterceptor()) // 在这里添加幂等拦截器 .transform(orderTransformer) .handle(orderService, "save") .get();
如果偏好更灵活的通道配置,也可以先定义一个带有拦截器的消息通道,再指定给消息驱动适配器的输出通道:
@Bean public MessageChannel jmsInputChannel(IdempotentReceiverInterceptor idempotentReceiverInterceptor) { DirectChannel channel = new DirectChannel(); channel.addInterceptor(idempotentReceiverInterceptor); return channel; } // 然后在适配器中绑定这个通道 .from(Jms.messageDrivenChannelAdapter(activeMQConnectionFactory) .destination(sourceQueue) .configureListenerContainer(spec -> { spec.sessionTransacted(false); spec.sessionAcknowledgeMode(Session.DUPS_OK_ACKNOWLEDGE); }) .outputChannel(jmsInputChannel())) .transform(orderTransformer) .handle(orderService, "save") .get();
第一种方式更简洁,推荐优先使用。
2. Oracle/MySQL中的元数据表结构
Spring Integration提供了JdbcMetadataStore来支持数据库存储幂等元数据,它依赖一张标准结构的表,MySQL和Oracle的表结构基本一致,仅需调整字符串类型细节:
MySQL表结构
CREATE TABLE INT_METADATA_STORE ( METADATA_KEY VARCHAR(255) NOT NULL PRIMARY KEY, METADATA_VALUE VARCHAR(255), REGION VARCHAR(100) NOT NULL DEFAULT 'DEFAULT' ); CREATE INDEX IDX_INT_METADATA_REGION ON INT_METADATA_STORE (REGION);
Oracle表结构
CREATE TABLE INT_METADATA_STORE ( METADATA_KEY VARCHAR2(255) NOT NULL PRIMARY KEY, METADATA_VALUE VARCHAR2(255), REGION VARCHAR2(100) NOT NULL DEFAULT 'DEFAULT' ); CREATE INDEX IDX_INT_METADATA_REGION ON INT_METADATA_STORE (REGION);
配置JdbcMetadataStore
接下来需要创建JdbcMetadataStore的Bean,并注入到幂等拦截器中:
@Bean public MetadataStore jdbcMetadataStore(DataSource dataSource) { JdbcMetadataStore metadataStore = new JdbcMetadataStore(dataSource); // 如果你的表名不是默认的INT_METADATA_STORE,可以在这里指定 // metadataStore.setTableName("YOUR_CUSTOM_TABLE_NAME"); // 可选:设置region隔离不同业务的幂等键,避免跨场景冲突 metadataStore.setRegion("ORDER_MESSAGE_REGION"); return metadataStore; } // 修改幂等拦截器,使用JDBC元数据存储 @Bean public IdempotentReceiverInterceptor idempotentReceiverInterceptor(MetadataStore jdbcMetadataStore) { IdempotentReceiverInterceptor idempotentReceiverInterceptor = new IdempotentReceiverInterceptor( new MetadataStoreSelector(m -> (String) m.getHeaders().get("JMSMessageId"), jdbcMetadataStore) ); idempotentReceiverInterceptor.setDiscardChannelName("ignoreDuplicates"); idempotentReceiverInterceptor.setThrowExceptionOnRejection(false); return idempotentReceiverInterceptor; }
REGION字段用于区分不同业务场景(比如不同的消息队列或业务模块),防止不同场景下的键值冲突。
内容的提问来源于stack exchange,提问作者nagendra
相关产品推荐
相关产品推荐

