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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:20:46