JavaEE定时任务需等待事务提交避免重复执行问题咨询
我有一些存储在MySQL数据库中的Task,Task主要用于调用远程系统的API,若两个Task并发执行会触发异常。当前实现代码如下:
@Singleton public class QueueScheduler { private final AtomicBoolean processing = new AtomicBoolean(false); @Inject public QueueService queueService; @Schedule(hour = "*", minute = "*", second = "*/5", persistent = false) @Lock(LockType.READ) public void tick() { if (processing.compareAndSet(false, true)) { try { queueService.executeNext(); } finally { processing.set(false); } } } }
@Dependent public class QueueService { @Inject RetryableService retryableService; @TransactionAttribute(TransactionAttributeType.REQUIRES_NEW) public void executeNext() { List<Queue> list = queueRepository.findLastIncomplete(1); if (list.isEmpty()) { return; } retryableService.tryExecute(task); } }
@Dependent public class RetryableService { public final EntityManager entityManager; public RetryableService(EntityManager entityManager) { this.entityManager = entityManager; } public <T extends RetryableEntity> void tryExecute(RetryableTask<T> task) { T entity = task.getEntity(); entity.setLastAttemptTimestamp(new Date()); entity.setAttemptCount(entity.getAttemptCount() + 1); try { task.run(entity); entity.setCompletedTimestamp(new Date()); entity.setCompletedSuccessfully(true); } catch (Exception e) { entity.setCompletedSuccessfully(false); entity.setCompletedTimestamp(new Date()); } entityManager.flush(); } }
@Stateless public class QueueRepository { @Inject public EntityManager entityManager; public List<Queue> findLastIncomplete(int count) { CriteriaBuilder cb = entityManager.getCriteriaBuilder(); CriteriaQuery<Queue> cq = cb.createQuery(Queue.class); Root<Queue> table = cq.from(Queue.class); cq.select(table); cq.where(cb.isNull(table.get(Queue_.completedTimestamp))); cq.orderBy(cb.desc(table.get(Queue_.startedTimestamp))); return entityManager.createQuery(cq) .setFirstResult(0) .setMaxResults(count) .getResultList(); } }
该实现多数情况运行正常,但偶尔会出现任务成功执行(无事务回滚、无报错)后,被另一个线程重复执行的情况。第二个线程在第一个线程完成后不到一秒内触发,读取到的Queue.completedTimestamp仍为null,从而重复执行任务。推测是QueueScheduler.tick()在实体完全持久化到数据库前触发了第二个线程(原本以为flush会完成持久化),导致读取到 stale 数据。
- 事务提交时机延迟:
RetryableService.tryExecute()里的entityManager.flush()仅将实体变更同步到数据库连接缓存,并未真正提交事务。事务提交是在QueueService.executeNext()方法结束时(因标注@TransactionAttribute(REQUIRES_NEW))触发,而QueueScheduler.tick()中processing.set(false)是在queueService.executeNext()返回后立即执行,此时事务可能仍在提交过程中,数据库未完成持久化,后续tick()线程就可能读到旧数据。 - Singleton锁类型错误:
QueueScheduler是@Singleton但使用@Lock(LockType.READ),意味着多个线程可同时进入tick()方法。虽用AtomicBoolean做控制,但READ锁本身允许并发调用,叠加事务提交延迟,极易出现竞态条件。
1. 修正Singleton的锁类型
将QueueScheduler.tick()的@Lock(LockType.READ)改为@Lock(LockType.WRITE),确保同一时间仅一个线程能进入tick()方法,从根源避免并发执行:
@Singleton public class QueueScheduler { @Inject public QueueService queueService; @Schedule(hour = "*", minute = "*", second = "*/5", persistent = false) @Lock(LockType.WRITE) public void tick() { queueService.executeNext(); } }
此时可移除AtomicBoolean processing,WRITE锁已保证方法串行执行。
2. 数据库层面加锁,避免读取 stale 数据
在QueueRepository.findLastIncomplete()中使用数据库行级锁,确保读取的任务未被其他线程处理:
public List<Queue> findLastIncomplete(int count) { CriteriaBuilder cb = entityManager.getCriteriaBuilder(); CriteriaQuery<Queue> cq = cb.createQuery(Queue.class); Root<Queue> table = cq.from(Queue.class); cq.select(table); cq.where(cb.isNull(table.get(Queue_.completedTimestamp))); cq.orderBy(cb.desc(table.get(Queue_.startedTimestamp))); // 使用悲观写锁,MySQL 8.0+可搭配SKIP LOCKED跳过已锁定行 return entityManager.createQuery(cq) .setFirstResult(0) .setMaxResults(count) .setLockMode(LockModeType.PESSIMISTIC_WRITE) .getResultList(); }
第一个线程读取任务时会加行锁,第二个线程会被阻塞或跳过已锁定行,直到第一个线程事务提交、锁释放后,才能读取到更新后的completedTimestamp。
3. 确保事务提交后再释放执行权限
若坚持使用AtomicBoolean,需确保事务完全提交后再设置processing为false,可通过JTA事务同步回调实现:
public <T extends RetryableEntity> void tryExecute(RetryableTask<T> task) { T entity = task.getEntity(); entity.setLastAttemptTimestamp(new Date()); entity.setAttemptCount(entity.getAttemptCount() + 1); try { task.run(entity); entity.setCompletedTimestamp(new Date()); entity.setCompletedSuccessfully(true); } catch (Exception e) { entity.setCompletedSuccessfully(false); entity.setCompletedTimestamp(new Date()); } // 注册事务同步回调,确保事务提交后再执行后续操作 TransactionSynchronizationRegistry tsr = CDI.current().select(TransactionSynchronizationRegistry.class).get(); tsr.registerInterposedSynchronization(new TransactionSynchronization() { @Override public void afterCompletion(int status) { if (status == Status.STATUS_COMMITTED) { // 事务提交完成后再释放执行权限 processing.set(false); } } }); }
优先选择修正Singleton锁类型为WRITE,确保调度方法串行执行;其次在数据库查询时添加悲观锁,从数据层面防止重复读取;最后通过事务同步回调确保事务提交完成后再允许下一次调度,彻底解决事务延迟带来的 stale 数据问题。
内容的提问来源于stack exchange,提问作者yoyo-san

