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

如何让JmsListenerContainerFactory的两个消费者顺序处理消息

JMS消费者并发访问同一条数据库记录的解决办法

现有配置

JmsListenerContainerFactory 配置

factory.setSessionTransacted(true);
factory.setConcurrency("1");
factory.setAutoStartup(true);

AtomikosConnectionFactoryBean 配置

atomikosfactory.setMaxPoolSize(2);

JMS消费者代码

@JmsListener(containerFactory = "jmsListenerContainerFactory", destination = "${queue-one}")
@Transactional
public <T> void handleOne(final Message<T> message)
{
   final var payload = message.getPayload();
   final var handler = getHandler(payload.getClass());
   handler.handleOne(payload);
}

@JmsListener(containerFactory = "jmsListenerContainerFactory", destination = "${queue-two}")
@Transactional
public <T> void handleTwo(final Message<T> message)
{
   final var payload = message.getPayload();
   final var handler = getHandler(payload.getClass());
   handler.handleTwo(payload);
}

问题

两个队列的消费者要操作数据库里的同一条记录,并发访问导致异常。试过让两个方法调用同一个@Synchronized标记的方法来串行执行,但出现事务回滚,且一些不能回滚的操作已经执行,后续线程重试时出错。

为啥@Synchronized没用

@Synchronized的同步范围是方法内部,但@Transactional的事务要等整个消费者方法执行完才提交。比如线程A进入同步块操作数据库,此时事务还没提交;线程B等着锁释放,等A的同步块跑完,A的事务可能因异常回滚,但如果A已经做了第三方接口调用、文件写入这类不可回滚的操作,B再执行就会因为数据不一致出问题。

解决方案

方案1:数据库加悲观锁

直接在查询共享记录时用SELECT ... FOR UPDATE(不同数据库语法可能略有差异),让数据库强制同一时间只有一个事务能获取该记录的锁,其他事务必须等待锁释放后才能操作。这样从数据库层面保证串行访问,避免冲突。

示例代码(在handler的数据库操作中):

// 查询时加悲观锁,其他事务需等待当前事务提交/回滚才能获取记录
SharedRecord record = entityManager.createQuery("SELECT r FROM SharedRecord r WHERE r.id = :id FOR UPDATE", SharedRecord.class)
    .setParameter("id", recordId)
    .getSingleResult();
// 后续正常操作record

方案2:给JMS容器配置单线程执行器

让两个消费者共用一个单线程的任务执行器,强制所有消息处理都在同一个线程中串行执行,从根源上杜绝并发。

修改JmsListenerContainerFactory配置:

// 创建单线程任务执行器
TaskExecutor singleThreadExecutor = new SimpleAsyncTaskExecutor();
singleThreadExecutor.setConcurrencyLimit(1);

factory.setSessionTransacted(true);
factory.setConcurrency("1");
factory.setAutoStartup(true);
// 为容器设置自定义任务执行器
factory.setTaskExecutor(singleThreadExecutor);

注意:这种方式会让两个队列的所有消息都串行处理,会降低整体吞吐量,适合对吞吐量要求不高但必须严格串行的场景。

方案3:使用细粒度本地锁(单实例应用适用)

如果是单实例部署,可针对共享记录的ID加细粒度锁,避免全局串行影响其他操作。确保锁覆盖整个数据库操作逻辑,同时控制锁的粒度。

示例代码:

// 用ConcurrentHashMap存储每个记录对应的锁对象,实现细粒度锁
private final ConcurrentHashMap<String, Object> recordLocks = new ConcurrentHashMap<>();

@JmsListener(containerFactory = "jmsListenerContainerFactory", destination = "${queue-one}")
@Transactional
public <T> void handleOne(final Message<T> message)
{
   final var payload = message.getPayload();
   // 从payload中获取共享记录的ID
   String targetRecordId = payload.getRecordId();
   // 获取对应记录的锁,不存在则创建
   Object lock = recordLocks.computeIfAbsent(targetRecordId, k -> new Object());
   
   synchronized (lock) {
       final var handler = getHandler(payload.getClass());
       handler.handleOne(payload);
   }
   // 清理锁,避免内存泄漏
   recordLocks.remove(targetRecordId);
}

// handleTwo方法执行同样的锁逻辑

这种方式仅对操作同一条记录的请求串行,其他请求可正常并发,对系统吞吐量影响较小。

注意事项

  • 避免使用全局锁(如锁整个类),会严重降低系统吞吐量;
  • 使用悲观锁时需注意数据库的锁超时时间,避免长时间占用锁导致其他请求阻塞;
  • 分布式部署场景下,本地锁无效,需使用Redis、Zookeeper等分布式锁组件,且要确保锁的释放与事务提交/回滚联动。

内容的提问来源于stack exchange,提问作者Effi T

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 11:53:12