Spring Batch并行处理时,如何实现单PDF多数据写入?
Spring Batch多线程下生成单份全量PDF解决方案
问题本质
启用SimpleTaskExecutor多线程批处理后,Spring Batch会将数据分片分配给多个独立工作线程,每个线程单独调用ItemWriter的write方法,因此无法在单次write调用中获取所有数据生成单份PDF。直接用ConcurrentHashMap或共享流的方式会因为线程隔离、资源竞争失效。
推荐方案:分阶段批处理(适配大数据量)
将任务拆分为两个独立Step,兼顾多线程处理效率和全量PDF生成需求:
阶段1:多线程并行处理并暂存数据
负责数据查询、校验,将处理后的结构化数据存入临时存储(数据库临时表/Redis/本地文件),保持多线程并行以提升效率。
临时存储Writer示例
@Component public class TempStorageWriter implements ItemWriter<PeopleDecoded> { @Autowired private TempPersonDecodedRepository tempRepo; @Override public void write(List<? extends PeopleDecoded> items) throws Exception { // 可在此预处理PDF所需字段,无需生成单个PDF tempRepo.saveAll(items); } }
阶段2:单线程生成全量PDF
读取临时存储中的所有数据,单线程生成包含全部内容的PDF,完成后可将PDF关联到原数据并清理临时存储。
全量数据Reader示例
@Component public class TempDataReader implements ItemReader<List<PeopleDecoded>> { @Autowired private TempPersonDecodedRepository tempRepo; private boolean readCompleted = false; @Override public List<PeopleDecoded> read() throws Exception { if (!readCompleted) { readCompleted = true; return tempRepo.findAll(); } return null; // 返回null表示读取完成 } }
全量PDF生成Writer示例
@Component public class FullPDFWriter implements ItemWriter<List<PeopleDecoded>> { @Autowired private PDFGeneratorService pdfService; @Autowired private PersonDecodedRepository finalRepo; @Autowired private TempPersonDecodedRepository tempRepo; @Override public void write(List<? extends List<PeopleDecoded>> items) throws Exception { List<PeopleDecoded> allData = items.get(0); File fullPDF = pdfService.createFullPDF(allData); // 将PDF字节关联到每条数据(按需调整) allData.forEach(item -> item.setFile(fullPDF.toBytes())); finalRepo.saveAll(allData); // 清理临时数据 tempRepo.deleteAll(); } }
Job配置示例
@Configuration public class BatchConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; @Bean public TaskExecutor taskExecutor() { SimpleTaskExecutor executor = new SimpleTaskExecutor(); executor.setConcurrencyLimit(20); return executor; } // 多线程处理Step @Bean public Step processAndStoreStep(ItemReader<PeopleDecoded> originalReader, ItemProcessor<PeopleDecoded, PeopleDecoded> processor, TempStorageWriter tempWriter) { return stepBuilderFactory.get("processAndStoreStep") .<PeopleDecoded, PeopleDecoded>chunk(100) // 合理设置Chunk大小 .reader(originalReader) .processor(processor) .writer(tempWriter) .taskExecutor(taskExecutor()) .build(); } // 单线程生成PDF Step @Bean public Step generateFullPDFStep(TempDataReader tempReader, FullPDFWriter fullPDFWriter) { return stepBuilderFactory.get("generateFullPDFStep") .<List<PeopleDecoded>, List<PeopleDecoded>>chunk(1) .reader(tempReader) .writer(fullPDFWriter) .build(); } @Bean public Job pdfGenerationJob() { return jobBuilderFactory.get("pdfGenerationJob") .start(processAndStoreStep(originalReader(), processor(), tempStorageWriter())) .next(generateFullPDFStep(tempDataReader(), fullPDFWriter())) .build(); } // 原有Reader、Processor Bean保持不变 @Bean public ItemReader<PeopleDecoded> originalReader() { // 原有PostgreSQL查询Reader配置 } @Bean public ItemProcessor<PeopleDecoded, PeopleDecoded> processor() { // 原有校验Processor配置 } }
备选方案:线程安全数据收集器(小数据量场景)
若数据量不大,可通过线程安全集合收集所有处理后数据,在Job结束时生成PDF。注意:大数据量会导致内存溢出,不推荐
数据收集器示例
@Component public class DataCollector { private final CopyOnWriteArrayList<PeopleDecoded> allData = new CopyOnWriteArrayList<>(); public void addItems(List<? extends PeopleDecoded> items) { allData.addAll(items); } public List<PeopleDecoded> getAllData() { return allData; } public void clear() { allData.clear(); } }
Job监听器(生成PDF)
@Component public class PDFGenerationListener implements JobExecutionListener { @Autowired private DataCollector dataCollector; @Autowired private PDFGeneratorService pdfService; @Autowired private PersonDecodedRepository repo; @Override public void beforeJob(JobExecution jobExecution) { dataCollector.clear(); } @Override public void afterJob(JobExecution jobExecution) { if (jobExecution.getStatus() == BatchStatus.COMPLETED) { List<PeopleDecoded> allData = dataCollector.getAllData(); File fullPDF = pdfService.createFullPDF(allData); allData.forEach(item -> item.setFile(fullPDF.toBytes())); repo.saveAll(allData); } } }
修改原有Writer
@Component public class CollectingWriter implements ItemWriter<PeopleDecoded> { @Autowired private DataCollector dataCollector; @Override public void write(List<? extends PeopleDecoded> items) throws Exception { // 预处理数据后加入收集器 dataCollector.addItems(items); } }
Job配置添加监听器
@Bean public Job pdfGenerationJob() { return jobBuilderFactory.get("pdfGenerationJob") .start(processingStep()) .listener(pdfGenerationListener()) .build(); }
关键注意事项
- 临时存储选择:分布式环境用数据库临时表/Redis,单机可考虑本地文件,但需确保任务失败后能清理。
- Chunk大小:多线程Step的Chunk不宜过小(避免频繁IO)或过大(内存压力),根据数据量调整。
- PDF生成优化:大数据量PDF需采用分页生成策略,避免内存溢出。
内容的提问来源于stack exchange,提问作者Ronak Joshi
相关产品推荐
相关产品推荐

