Spring多实例批处理:如何避免重复获取待处理记录?
我之前帮别人处理过类似的多实例Spring Batch并发问题,你的核心痛点其实是跨JVM进程的记录抢占——单进程的synchronized只能管自己进程内的线程,普通的悲观锁又会让后续事务等待,等锁释放后还是可能查到已被更新的记录,导致无效操作。结合你不能“先更新再查询”的要求,给你几个落地性强的方案:
方案1:数据库行级锁 + 跳过已锁定记录(优先推荐)
这个方案不需要额外中间件,直接利用数据库的原生锁特性解决问题,核心是用SELECT FOR UPDATE SKIP LOCKED语法(MySQL 8.0+、PostgreSQL、Oracle等主流数据库都支持)。
原理
当事务1执行带SKIP LOCKED的查询时,会锁定查询到的100条记录,事务2的相同查询会直接跳过这些已锁定的记录,去获取下一批未被处理的NOT COMPLETED记录,完全避免了重复获取的情况。
代码改造
首先修改查询方法的SQL,加上锁和跳过逻辑:
// DAO层方法 @Query(value = "SELECT * FROM student WHERE status = 'NOT COMPLETED' ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED", nativeQuery = true) List<Student> getNotCompletedRecordsWithSkipLocked();
然后调整业务方法,去掉没用的synchronized(多实例下它不起作用),同时给update语句加双重保险:
@Transactional(isolation = Isolation.READ_COMMITTED) public List<Student> processStudentRecords() { List<Student> students = getNotCompletedRecordsWithSkipLocked(); if (!students.isEmpty()) { // 提取记录ID,更新时额外加status条件,防止极端情况的重复更新 List<Long> studentIds = students.stream().map(Student::getId).collect(Collectors.toList()); updateStatusToInProgress(studentIds); } return students; } // DAO层更新方法 @Modifying @Query(value = "UPDATE student SET status = 'IN PROGRESS' WHERE id IN (:ids) AND status = 'NOT COMPLETED'", nativeQuery = true) void updateStatusToInProgress(@Param("ids") List<Long> ids);
注意事项
- 事务隔离级别用
READ_COMMITTED即可,避免不必要的锁持有时间; - 如果你的数据库不支持
SKIP LOCKED(比如老版本MySQL),可以用SELECT FOR UPDATE NOWAIT,但它会在遇到锁时直接抛出异常,需要你捕获异常并重试查询逻辑。
方案2:Spring Batch分区机制
如果你的数据量较大(1000万条),可以用Spring Batch的**分区(Partitioning)**功能,把数据分成多个独立的分区,每个实例只处理自己分区内的记录,从根源上避免并发冲突。
原理
比如按记录的ID范围或者哈希值分区,比如把ID分成5个区间,每个实例处理一个区间的记录,多实例之间不会交叉访问同一段数据,自然不会出现重复获取的问题。
示例配置
@Bean public Step partitionStep(StepBuilderFactory stepBuilderFactory, Partitioner partitioner, Step slaveStep) { return stepBuilderFactory.get("studentPartitionStep") .partitioner("slaveStep", partitioner) .step(slaveStep) .gridSize(5) // 分区数量,对应你的实例数或线程数 .taskExecutor(new SimpleAsyncTaskExecutor()) .build(); } @Bean public Partitioner studentPartitioner(DataSource dataSource) { ColumnRangePartitioner partitioner = new ColumnRangePartitioner(); partitioner.setColumn("id"); // 按ID分区 partitioner.setDataSource(dataSource); partitioner.setTable("student"); return partitioner; }
之后你的业务方法只需要处理当前分区内的NOT COMPLETED记录即可,完全不需要锁机制。
方案3:分布式锁控制查询批次
如果以上两种方案都不适用,你可以用Redis或ZooKeeper实现分布式锁,确保同一时间只有一个实例的一个线程能执行“查询+更新”操作。
原理
在进入processStudentRecords()方法前,先获取一个全局分布式锁,获取到锁的线程才能执行查询和更新操作,其他线程/实例需要等待锁释放后再尝试。不过这个方案会降低吞吐量,适合数据处理优先级不高的场景。
代码示例(Redis锁为例)
@Autowired private RedisLockService redisLockService; @Transactional(isolation = Isolation.READ_COMMITTED) public List<Student> processStudentRecords() { // 获取分布式锁,超时时间设为事务最长执行时间 RedisLock lock = redisLockService.acquireLock("student-processing-lock", 30, TimeUnit.SECONDS); if (lock == null) { // 没获取到锁,返回空列表或重试 return Collections.emptyList(); } try { List<Student> students = getNotCompletedRecords(); if (!students.isEmpty()) { updateStatusToInProgress(students); } return students; } finally { // 释放锁 lock.release(); } }
为什么之前的方案没用?
你提到用了乐观锁、悲观锁还是重复获取,核心原因是:
- 普通的
SELECT FOR UPDATE会让事务2等待,等事务1提交后,事务2的查询会返回已经被更新为IN PROGRESS的记录,这时候你的update操作如果没加status = 'NOT COMPLETED'的条件,就会执行无效的更新; - 乐观锁是基于版本号的,需要每次更新时校验版本,适合低并发场景,高并发下会有大量更新失败的情况,不适合你的1000万条数据场景。
内容的提问来源于stack exchange,提问作者Tech Guy

