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

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 数据。


问题根源分析
  1. 事务提交时机延迟:RetryableService.tryExecute()里的entityManager.flush()仅将实体变更同步到数据库连接缓存,并未真正提交事务。事务提交是在QueueService.executeNext()方法结束时(因标注@TransactionAttribute(REQUIRES_NEW))触发,而QueueScheduler.tick()中processing.set(false)是在queueService.executeNext()返回后立即执行,此时事务可能仍在提交过程中,数据库未完成持久化,后续tick()线程就可能读到旧数据。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 03:37:21