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

Spring多实例批处理:如何避免重复获取待处理记录?

解决多实例Spring Batch并发获取记录的问题

我之前帮别人处理过类似的多实例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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:59:22