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()); } }
问题分析与解决方案
核心问题:分页读取与多线程更新的竞态条件
- 无行级锁定导致并发冲突:多线程环境下,多个线程可能同时读取到同一条待处理记录,或者某条记录被一个线程处理但未提交时,另一个线程的分页查询仍会包含它,最终导致部分记录被跳过或重复处理。
- 手动事务与Spring Batch事务冲突:你在
batchOperation中手动管理事务,会和Spring Batch为每个chunk自动创建的事务边界冲突,导致部分更新未正确提交或回滚。 - 分页读取的快照特性: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
相关产品推荐
相关产品推荐

