Spring Batch多线程步骤中TaskExecutor线程重复执行Reader问题求助
解决Spring Batch多线程分块处理中重复读取数据的问题
问题场景
开发从BigQuery获取数据的Spring Batch作业,设计逻辑如下:
TripDateTimeDecider:确定Reader执行查询的时间范围TransactionReader:执行BigQuery查询加载数据TransactionProcessor:处理加载的数据TransactionWriter:将处理后的数据写入目标表
预期流程:TripDateTimeDecider → TransactionReader(单次读取数据)→ 分块数据多线程执行Processor和Writer
实际问题:多个TaskExecutor线程重复调用TransactionReader读取相同时间范围的数据,导致Processor和Writer重复处理同一批数据
运行日志片段
2023-04-11 12:50:57.456 [taskExecutor-3] INFO c.q.p.p.steps.TransactionReader - TransactionReader::read() for tripStartDateTime= 2022-03-01T00:00:00 and tripIntervalDateTime= 2022-03-01T06:00:00.0 2023-04-11 12:51:01.286 [taskExecutor-3] INFO c.q.p.p.utils.BigQuerySalesTransUtil - loadTransactionsFromURT for trip_start_date_time=2022-03-01T00:00:00 , tripIntervalDateTime= 2022-03-01T06:00:00.0 and currentEnv = dev 2023-04-11 12:51:01.287 [taskExecutor-4] INFO c.q.p.p.steps.TransactionReader - TransactionReader::read() for tripStartDateTime= 2022-03-01T00:00:00 and tripIntervalDateTime= 2022-03-01T06:00:00.0 2023-04-11 12:51:01.287 [taskExecutor-4] INFO c.q.p.p.utils.BigQuerySalesTransUtil - loadTransactionsFromURT for trip_start_date_time=2022-03-01T00:00:00 , tripIntervalDateTime= 2022-03-01T06:00:00.0 and currentEnv = dev
当前作业核心配置
@Bean protected Step processLines() { return steps.get("processEntities").<TransactionReceiptScanRequest, TransactionReceiptScanRequest> chunk(10) .reader(transactionReader(WILL_BE_INJECTED,WILL_BE_INJECTED,WILL_BE_INJECTED,WILL_BE_INJECTED)) .processor(transactionProcessor()) .writer(transactionWriter()) .taskExecutor(taskExecutor()) .build(); }
问题根源
直接给Chunk步骤配置taskExecutor时,Spring Batch会启动多个独立线程,每个线程会完整执行读取→处理→写入的全流程。这导致每个线程都会调用TransactionReader的read()方法,重复读取相同的数据源。
解决方案
使用Spring Batch提供的AsyncItemProcessor和AsyncItemWriter组件,实现单线程读取数据,多线程处理/写入分块数据的目标。
步骤1:配置异步Processor和Writer
新增Bean配置,用异步组件包装原有的Processor和Writer:
@Bean @StepScope public AsyncItemProcessor<TransactionReceiptScanRequest, TransactionReceiptScanRequest> asyncTransactionProcessor(TaskExecutor taskExecutor) { AsyncItemProcessor<TransactionReceiptScanRequest, TransactionReceiptScanRequest> asyncProcessor = new AsyncItemProcessor<>(); asyncProcessor.setDelegate(transactionProcessor()); // 包装原有业务Processor asyncProcessor.setTaskExecutor(taskExecutor); // 指定处理用线程池 return asyncProcessor; } @Bean @StepScope public AsyncItemWriter<TransactionReceiptScanRequest> asyncTransactionWriter() { AsyncItemWriter<TransactionReceiptScanRequest> asyncWriter = new AsyncItemWriter<>(); asyncWriter.setDelegate(transactionWriter()); // 包装原有业务Writer return asyncWriter; }
步骤2:修改Chunk步骤配置
更新processLines步骤,替换为异步组件,移除直接给Chunk配置的taskExecutor:
@Bean protected Step processLines() { return steps.get("processEntities") // 注意泛型变更:输出类型为Future<TransactionReceiptScanRequest> .<TransactionReceiptScanRequest, Future<TransactionReceiptScanRequest>> chunk(10) .reader(transactionReader(WILL_BE_INJECTED,WILL_BE_INJECTED,WILL_BE_INJECTED,WILL_BE_INJECTED)) .processor(asyncTransactionProcessor(taskExecutor())) .writer(asyncTransactionWriter()) .build(); }
原理说明
AsyncItemProcessor:主线程读取数据后,将每个Item的处理任务提交到线程池异步执行,返回Future对象AsyncItemWriter:等待所有异步处理任务完成,收集结果后调用原Writer完成写入- Reader始终在主线程单线程执行,彻底避免重复读取数据
内容的提问来源于stack exchange,提问作者NIRAJ KUMAR
相关产品推荐
相关产品推荐

