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

