如何在Groovy中解决并发偏差?队列多线程竞态问题咨询
哈哈,这个问题我太熟了!你遇到的是典型的竞态条件(Race Condition)——查询状态为wait的队列,和把它们更新为running这两个操作不是原子性的。当多个线程同时执行时,很容易出现线程A刚查到一批wait的队列,还没来得及更新状态,线程B也查到了同一批队列,最后两个线程都去处理同一个队列的尴尬情况。
下面给你几个Java生态里常用的解决思路,都是实战中验证过的:
1. 数据库层面的原子更新(最推荐)
如果你的队列数据存在数据库里,直接把查询和更新合并成原子性的数据库操作,这是最直接有效的方案,数据库会帮你保证操作的排他性。
方案1:直接原子更新
用JPQL或者原生SQL写一个更新语句,一步完成状态变更:
@Repository public interface QueueRepository extends JpaRepository<Queue, Long> { @Modifying @Query("UPDATE Queue q SET q.state = 'running' WHERE q.state = 'wait'") int updateWaitToRunning(); }
这个方法会返回被成功更新的行数,如果你需要拿到具体的队列对象,可以用数据库的行锁语法,比如MySQL 8+、PostgreSQL都支持的SELECT ... FOR UPDATE SKIP LOCKED:
@Repository public interface QueueRepository extends JpaRepository<Queue, Long> { // 锁定符合条件的行,其他线程会跳过这些被锁定的行 @Query(value = "SELECT * FROM queue WHERE state = 'wait' FOR UPDATE SKIP LOCKED", nativeQuery = true) List<Queue> findAndLockWaitQueues(); } // 然后在事务中执行查询+更新 @Service public class QueueService { @Autowired private QueueRepository queueRepository; @Transactional public List<Queue> getAndStartQueues() { List<Queue> queues = queueRepository.findAndLockWaitQueues(); queues.forEach(queue -> queue.setState("running")); queueRepository.saveAll(queues); return queues; } }
FOR UPDATE SKIP LOCKED的作用是:查询时锁定wait状态的队列,其他线程查询时会自动跳过这些被锁定的行,从根源上避免了重复获取队列的问题。
2. 乐观锁机制(适合低并发场景)
如果不想用数据库的行锁,可以给Queue实体加一个版本号字段,利用JPA的乐观锁来避免并发更新冲突:
@Entity public class Queue { // 其他字段... private String state; // 乐观锁版本号,JPA会自动管理这个字段的更新 @Version private Integer version; }
更新时,JPA会自动检查版本号:如果当前线程拿到的版本号和数据库里的不一致,说明这个队列已经被其他线程更新过了,会抛出OptimisticLockingFailureException,这时候你可以捕获异常并重试:
@Service public class QueueService { @Autowired private QueueRepository queueRepository; @Transactional public void updateQueueState(Long queueId) { boolean updateSuccess = false; while (!updateSuccess) { try { Queue queue = queueRepository.findById(queueId).orElseThrow(() -> new RuntimeException("队列不存在")); if ("wait".equals(queue.getState())) { queue.setState("running"); queueRepository.save(queue); } updateSuccess = true; } catch (OptimisticLockingFailureException e) { // 版本不匹配,说明被其他线程抢先更新了,短暂等待后重试 try { Thread.sleep(100); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); break; } } } } }
这种方式不需要加数据库锁,性能损耗小,但需要处理重试逻辑,适合并发量不是特别高的场景。
3. 分布式锁(多实例部署场景)
如果你的应用是多实例部署的,上面的数据库锁可能覆盖不到跨实例的线程竞争,这时候可以用分布式锁,比如基于Redis的Redisson:
@Service public class QueueService { @Autowired private RedissonClient redissonClient; @Autowired private QueueRepository queueRepository; public void processQueues() { // 定义锁的名称,保证全局唯一 RLock lock = redissonClient.getLock("queue-processing-lock"); try { // 尝试获取锁:最多等待10秒,持有锁1分钟 if (lock.tryLock(10, 60, TimeUnit.SECONDS)) { // 拿到锁后再执行查询和更新逻辑,保证同一时间只有一个线程(跨实例)执行 List<Queue> queues = queueRepository.findAllByState("wait"); queues.forEach(queue -> queue.setState("running")); queueRepository.saveAll(queues); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { // 记得释放锁,避免死锁 if (lock.isHeldByCurrentThread()) { lock.unlock(); } } } }
如果担心锁粒度太粗影响性能,可以考虑按队列的分组ID或者范围来拆分锁,减少锁竞争。
内容的提问来源于stack exchange,提问作者Dasma

