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

Spring Batch如何实现Reader、Processor、Writer多线程处理?

Spring Batch批量文件处理性能优化问题

需求说明

我正在使用Spring Batch应用进行批量文件处理,流程如下:

  • 通过网络调用读取文件;
  • 准备批量JSON请求体并调用接口;
  • 将响应写入文件。

现存问题

流程功能正常但性能极低,处理含25K条记录的小文件通常需要16分钟,原因如下:

  • Reader调用Processor时会阻塞等待;
  • 响应就绪后,写入操作因IO缓慢再次阻塞。

耗时假设:

  • 读取并准备JSON:2秒 [READER]
  • 接口请求:2秒 [Processor]
  • 写入文件:1秒 [Writer]

单线程阻塞调用流程:

Single Threaded Block Calls: READER --> Processor --> Writer. //Total 5 seconds per request.

期望处理模式

多线程阻塞调用:

|- - > Processor - -> Writer.
READER -- >|- - > Processor - -> Writer.
          |- - > Processor - -> Writer.

当前配置代码

@Bean
  public PlatformTransactionManager transactionManager() {
    return new JpaTransactionManager();
  }

  @Bean
  @Autowired
  Step step1(JobRepository jobRepository) {
    PlatformTransactionManager transactionManager = transactionManager();
    StepBuilder stepBuilder = new StepBuilder("CARD TRANSFORMATION", jobRepository);
    return stepBuilder
        .<List<FileStructure>, CardExtractOutputList>chunk(1, transactionManager)
        .reader(generalFileReader.reader(""))
        .processor(cardExtractFileProcessor)
        .writer(cardExtractFileWriter)
        .taskExecutor(taskExecutor())
        .faultTolerant()
        .retryLimit(3)
        .retry(RuntimeException.class)
        .build();
  }



  @Bean(name = "jsob")
  @Autowired
  Job cardExtractFilejob(JobRepository jobRepository) {
    JobBuilder jobBuilderFactory =
        new JobBuilder("somename", jobRepository)
            .incrementer(new RunIdIncrementer())
            .listener(preJobListener)
            .listener(postJobExecution);
    return jobBuilderFactory.flow(step1(jobRepository)).end().build();
  }


  @Bean
  public TaskExecutor taskExecutor() {
    SimpleAsyncTaskExecutor asyncTaskExecutor = new SimpleAsyncTaskExecutor();
    asyncTaskExecutor.setConcurrencyLimit(10);
    return asyncTaskExecutor;
  }

自定义Reader代码

@Bean
  @StepScope
  @SneakyThrows
  public MultiLinePeekableReader reader(
      @Value(FILENAME_JOB_PARAM) final String fileName) {
    FlatFileItemReader<FileStructure> itemReader = new FlatFileItemReader<>() {};
    final String gcsLocationOfFile =
        FilesUtility.getAbsoluteGCSPathOfAFile(fileName, gcsRelatedConfiguration);

    final Resource resource = applicationContext.getResource(gcsLocationOfFile);
    itemReader.setResource(resource);
    itemReader.setName("FileReader : " + fileName);
    itemReader.setLineMapper(lineMapper());
    itemReader.setStrict(true);

    MultiLinePeekableReader multiLinePeekableReader = new MultiLinePeekableReader(fileName);
    multiLinePeekableReader.setDelegate(itemReader);
    return multiLinePeekableReader;
  }

  private LineMapper<FileStructure> lineMapper() {
    DefaultLineMapper<FileStructure> lineMapper = new DefaultLineMapper<>();
    DelimitedLineTokenizer lineTokenizer = new DelimitedLineTokenizer();
      ..
    return lineMapper;
  }
}

MultiLinePeekableReader代码

public class MultiLinePeekableReader implements ItemReader<List<FileStructure>>, ItemStream {
 private SingleItemPeekableItemReader<FileStructure> delegate;
..
  @Override
  @SneakyThrows
  @Bean
  @StepScope
  public synchronized List<FileStructure> read() {
    List<FileStructure> records = null;
    int readCount = fileProcessingConfiguration.itemsPerRead();
    try {
      for (FileStructure line; (line = this.delegate.read()) != null; ) {
        seqNo = seqNo.add(new BigInteger(FileProcessingConstants.NUMBER_STRING_ONE));
        line.setSequenceNo(seqNo.toString());
        line.setMaskedSensitiveData(
            FilesUtility.getMaskedSensitiveDataFromData(
                line.getSensitiveData(),
                fileProcessingConfiguration.leadingPersistCount(),
                fileProcessingConfiguration.trailingPersistCount()));

        if (readCount == fileProcessingConfiguration.itemsPerRead()) {
          records = new ArrayList<>();
          records.add(line);
          readCount--;
        } else {
          records.add(line);
          readCount--;
          FileStructure nextLine = this.delegate.peek();
          if (nextLine == null || readCount == 0) {
            readCount = fileProcessingConfiguration.itemsPerRead();
            return records;
          }
        }
      }
    } catch (FlatFileParseException parseException) {
      if (records == null) {
        records = new ArrayList<>();
      }
      ..
    } 
    return records;
  }

  @Override
  public void close() throws ItemStreamException {
    this.delegate.close();
  }

  @Override
  public void open(ExecutionContext executionContext) throws ItemStreamException {
    this.delegate.open(executionContext);
  }

  @Override
  public void update(ExecutionContext executionContext) throws ItemStreamException {
    this.delegate.update(executionContext);
  }

  public void setDelegate(FlatFileItemReader<FileStructure> delegate) {
    this.delegate = new SingleItemPeekableItemReader<>();
    this.delegate.setDelegate(delegate);
  }

}

已查阅无帮助的资料

  • Spring Batch单线程Reader多线程Writer相关内容

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 23:45:17