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
相关产品推荐
相关产品推荐

