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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 08:15:04