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

Spring Batch实现Process Indicator Pattern遇未处理记录问题求助

问题:使用JdbcPagingItemReader多线程处理时部分记录未被处理

我正在寻找Process Indicator Pattern的示例。当前我有一个源表STUDENTS,其中STATUS字段用于标识记录是否已处理,采用任务执行器进行多线程处理。

在Writer中,我会将处理后的记录插入新表PROCESSED_STUDENTS,同时把源表STUDENTS中已处理记录的STATUS更新为Processed,这两个操作放在事务块中执行,确保故障时能回滚更改。

但结合JdbcPagingItemReader使用后,流程结束仍有部分记录未被处理,请问我遗漏了什么?


Reader 代码

@Bean
@StepScope
public ItemReader<SourceData> reader(DataSource dataSource) {
        
    Map<String, Object> parameterValues = new HashMap<>();
    parameterValues.put("status", "ToBeProcessed");     
    
    JdbcPagingItemReader<SourceData> reader = new JdbcPagingItemReader<>();
    reader.setName("Oracle_RCP");
    reader.setDataSource(dataSource);
    reader.setRowMapper(SourceData.rowMapper());        
    reader.setParameterValues(parameterValues);
    reader.setPageSize(100);        
    reader.setQueryProvider(getQueryProvider(new OraclePagingQueryProvider(), "SELECT ID, NAME, CREATED_TIME", "FROM STUDENTS", "WHERE STATUS = :status", CREATED_TIME, Order.ASCENDING));
    reader.setSaveState(false);     
    
    try {
        reader.afterPropertiesSet();
    } catch (Exception e) {
        log.error(e.getMessage(), e.getStackTrace());
    }
    
    return reader;
}

public PagingQueryProvider getQueryProvider(AbstractSqlPagingQueryProvider queryProvider, String select, String from, String where, String sortKey, Order order) {
    queryProvider.setSelectClause(select);
    queryProvider.setFromClause(from);
    
    if (where != null) {
        queryProvider.setWhereClause(where);
    }
    
    Map<String, Order> sortConfiguration = new HashMap<>();
    sortConfiguration.put(sortKey, order);
    queryProvider.setSortKeys(sortConfiguration);
    
    return queryProvider;
}

Processor 代码

@Bean
@StepScope
public ItemProcessor<SourceData, OutData> processor(
        @Value("#{jobParameters['processDate']}") String processDate) {
    return new CustomItemProcessor(processDate);        
}

Writer 代码

@SuppressWarnings("unchecked")
@Bean
public ItemWriter<OutData> writer(Utils utils) {        
    return OutDataList -> utils.batchOperation((List<OutData>) OutDataList, chunk); 
}

Job、Step 及 TaskExecutor 代码

@Bean("MainJob")
@Scope(value = ConfigurableBeanFactory.SCOPE_PROTOTYPE)
public Job mainJob(JobBuilderFactory jobBuilderFactory, Step step) {        
    return jobBuilderFactory.get("mainJob")   
      .incrementer(new RunIdIncrementer())          
      .flow(step)
      .end()
      .build();
}

@Bean
public Step step(StepBuilderFactory stepBuilderFactory) {        
    return stepBuilderFactory.get("step")               
            .<SourceData, OutData> chunk(chunk)
            .reader(reader(null))
            .processor(processor(null))
            .writer(writer(null))   
            .taskExecutor(taskExecutor())
            .build();
}

@Bean
@StepScope
public ThreadPoolTaskExecutor taskExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setMaxPoolSize(50);
    executor.setCorePoolSize(25);       
    executor.setQueueCapacity(25);
    executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
    executor.setThreadNamePrefix("MultiThreaded-Executor");
    return executor;
}

Writer 中的 batchOperation 方法

public void batchOperation(List<OutData> outDataList, int batchSize) {
    try (
            Connection con = jdbcTemplate.getDataSource().getConnection(); 
            PreparedStatement psInsert = con.prepareStatement("INSERT INTO PROCESSED_STUDENTS (ID, NAME) VALUES (?, ?)");
            PreparedStatement psUpdate = con.prepareStatement("UPDATE STUDENTS SET STATUS = 'Processed' WHERE ID = ?");) {
        
        // Starting transaction block
        con.setAutoCommit(false);
        
        int i = 0;
        for (OutData argument : outDataList) {          
            psInsert.setLong(1, argument.getId());
            psInsert.setString(2, argument.getName());
                        
            psUpdate.setLong(1, argument.getId());
                        
            psInsert.addBatch();
            psUpdate.addBatch();
                        
            i++;
                        
            if (i % batchSize == 0) {
                psInsert.executeBatch();
                psUpdate.executeBatch();
            }               
        }
        
        // Executing remaining batch if total record count is an odd number
        psInsert.executeBatch();
        psUpdate.executeBatch();
        
        // End transaction block, commit changes
        con.commit();

        // Setting it back to default true
        con.setAutoCommit(true);
        
    } catch (Exception e) {
        log.error(e.getMessage());
    }
}

问题分析与解决方案

核心问题:分页读取与多线程更新的竞态条件

  1. 无行级锁定导致并发冲突:多线程环境下,多个线程可能同时读取到同一条待处理记录,或者某条记录被一个线程处理但未提交时,另一个线程的分页查询仍会包含它,最终导致部分记录被跳过或重复处理。
  2. 手动事务与Spring Batch事务冲突:你在batchOperation中手动管理事务,会和Spring Batch为每个chunk自动创建的事务边界冲突,导致部分更新未正确提交或回滚。
  3. 分页读取的快照特性:Oracle的分页查询基于当前提交的数据,当线程更新记录状态并提交后,后续分页查询会排除这些记录,但如果初始分页已加载未提交的记录,可能导致处理逻辑混乱。

具体修复步骤

1. 读取时添加行级锁定

修改Reader的查询条件,添加Oracle支持的FOR UPDATE SKIP LOCKED,读取时锁定待处理记录,其他线程自动跳过已锁定的记录,确保每个线程处理唯一的记录:

// 修改getQueryProvider中的where子句
queryProvider.setWhereClause("WHERE STATUS = :status FOR UPDATE SKIP LOCKED");

2. 移除手动事务管理,依赖Spring Batch事务

删除batchOperation中的手动事务代码,由Spring Batch统一管理chunk级事务,确保插入和更新操作的原子性:

public void batchOperation(List<OutData> outDataList, int batchSize) {
    try (
            Connection con = jdbcTemplate.getDataSource().getConnection(); 
            PreparedStatement psInsert = con.prepareStatement("INSERT INTO PROCESSED_STUDENTS (ID, NAME) VALUES (?, ?)");
            PreparedStatement psUpdate = con.prepareStatement("UPDATE STUDENTS SET STATUS = 'Processed' WHERE ID = ?");) {
        
        int i = 0;
        for (OutData argument : outDataList) {          
            psInsert.setLong(1, argument.getId());
            psInsert.setString(2, argument.getName());
                        
            psUpdate.setLong(1, argument.getId());
                        
            psInsert.addBatch();
            psUpdate.addBatch();
                        
            i++;
                        
            if (i % batchSize == 0) {
                psInsert.executeBatch();
                psUpdate.executeBatch();
            }               
        }
        
        psInsert.executeBatch();
        psUpdate.executeBatch();
        
    } catch (Exception e) {
        log.error(e.getMessage());
        throw new RuntimeException(e); // 抛出异常触发Spring Batch事务回滚
    }
}

3. 调整TaskExecutor配置

增大队列容量并手动初始化执行器,避免任务溢出导致的记录遗漏:

@Bean
@StepScope
public ThreadPoolTaskExecutor taskExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setMaxPoolSize(50);
    executor.setCorePoolSize(25);       
    executor.setQueueCapacity(100); // 增大队列容量
    executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
    executor.setThreadNamePrefix("MultiThreaded-Executor");
    executor.initialize(); // 手动初始化确保执行器正常启动
    return executor;
}

4. 优化Process Indicator逻辑(可选)

可以将状态字段扩展为ToBeProcessed、Processing、Processed、Failed,读取时将状态改为Processing,处理完成后改为Processed,失败时改为Failed,进一步避免并发冲突和重复处理。


内容的提问来源于stack exchange,提问作者Somebody

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 22:37:12