Spring Batch处理2-3万条数据耗时超1小时,求性能优化方案
Spring Batch处理2-3万条数据耗时超1小时的优化方案
问题背景
使用Spring Batch处理2-3万条数据时耗时超过1小时,当前技术栈与配置:
- 核心组件:ItemStreamReader、ItemProcessor、JpaItemWriter
- 运行模式:分区模式,配置grid size=6、chunk size=500、batch-size=1000
- 数据库:H2
当前配置与代码
分区Step配置类
@Configuration @ConditionalOnAderBatchStep(BATCH_NAME) @Import({AsynchronousTasksConfiguration.class, BankAccountsReaderWriterConfiguration.class}) @Slf4j public class ParallelClientBankAccountImportStepConfiguration { public static final String STEP_NAME = "parallelClientBankAccountsImportStep"; public static final String ITEM_WRITER_NAME = "bankAccountsItemWriter"; public static final String ITEM_READER_NAME = "bankAccountsItemReader"; @Bean public Step parallelClientBankAccountsImportStep(StepBuilderFactory stepBuilderFactory, PlatformTransactionManager jpaTransactionManager, @Value("${batch.export.parallel.grid-size:6}") int gridSize, ItemWriter<BankAccount> itemWriter, TaskExecutor applicationTaskExecutor, ApplicationEventPublisher applicationEventPublisher) { log.info("Configuring the [{}].", STEP_NAME); String slaveStepName = "slaveBankAccountsImportStep"; return stepBuilderFactory.get(STEP_NAME) .partitioner(slaveStepName, partitioner()) .gridSize(gridSize) .step(slaveStep(slaveStepName, stepBuilderFactory, itemWriter, jpaTransactionManager, applicationEventPublisher)) .taskExecutor(applicationTaskExecutor) .transactionManager(jpaTransactionManager) .build(); } protected final Partitioner partitioner() { return gridSize -> { Map<String, ExecutionContext> map = new HashMap<>(gridSize); for (int i = 1; i <= gridSize; i++) { ExecutionContext context = new ExecutionContext(); context.put("index", i); map.put("partition_" + i, context); } return map; }; } protected Step slaveStep(String slaveStepName, StepBuilderFactory stepBuilderFactory, ItemWriter<BankAccount> itemWriter, PlatformTransactionManager jpaTransactionManager, ApplicationEventPublisher applicationEventPublisher) { return stepBuilderFactory.get(slaveStepName) .transactionManager(jpaTransactionManager) .<ClientBankAccount, BankAccount>chunk(500) .reader(itemReader(null, 0, 0)) .processor(itemProcessor(null)) .writer(itemWriter) .listener(chunkLogger()) .listener(stepLogger(applicationEventPublisher)) .build(); } @StepScope @Bean public ItemStreamReader<ClientBankAccount> itemReader(ClientBankAccountsItemReader itemReader, @Value("${batch.export.parallel.grid-size:6}") int gridSize, @Value("#{stepExecutionContext[index]}") int partitionIndex) { ItemStreamReader<ClientBankAccount> itemStreamReader = new ClientParallelItemStreamReader(itemReader, gridSize, partitionIndex); return itemStreamReader; } @Bean @StepScope public ItemProcessor<ClientBankAccount, BankAccount> itemProcessor(ClientBankAccountMapper mapper) { ItemProcessor<ClientBankAccount, BankAccount> itemProcessor = mapper::fromClientService; // 注:原代码中startTime未定义,需移除或修正 // long elapsedTime = startTime - System.currentTimeMillis(); return itemProcessor; } }
ItemReader实现类
@Slf4j @Service public class ClientBankAccountsItemReader extends AbstractItemCountingItemStreamItemReader<ClientBankAccount> { private final ReactiveClientService clientService; private Disposable clientBankAccounts; private BlockingQueue<ClientBankAccount> queue = new LinkedBlockingQueue<>(); private boolean isOpen = false; public ClientBankAccountsItemReader(ReactiveClientService clientService) { this.clientService = clientService; setName("ClientBankAccountsItemReader"); } public synchronized boolean isOpen() { return this.isOpen; } @Override protected synchronized void doOpen()throws Exception { if (!isOpen) { log.info("Opening the Client Account Reader..."); clientBankAccounts = clientService.getAccountsApi(getSearchAccountRequest()).subscribe(this::addToQueue); this.isOpen = true; } } private void addToQueue(ClientBankAccount ba) { try { if (!queue.offer(ba, 5, MILLISECONDS)) { log.warn("Could not add Client Bank Account to the internal queue!"); log.trace("Client Bank Account: {}", ba); } } catch (InterruptedException e) { log.warn("INTERRUPTED! Could not ADD Client Bank Account to the internal queue!"); } } @Override protected void doClose() { if (this.clientBankAccounts != null) { this.clientBankAccounts.dispose(); } } @Override protected ClientBankAccount doRead() { try { return this.queue.poll(50, SECONDS); } catch (InterruptedException e) { log.warn("INTERRUPTED! Could not READ client Bank Account from the internal queue!"); } return null; } private SearchAccountRequest getSearchAccountRequest(){ SearchAccountRequest searchAccountRequest= new SearchAccountRequest(); // request details are set return searchAccountRequest; } }
ItemWriter实现类
@Slf4j public class BankAccountJpaItemWriter extends JpaItemWriter<BankAccount> { private final JpaItemWriter<CABankAccount> closedAccountsWriter; // declared few repository... public BankAccountJpaItemWriter(Parms...) { // parameterised constructor } @Override public void write(List<? extends BankAccount> items) { if (isEmpty(items)) { return; } List<BankAccount> activeBankAccounts = new ArrayList<>(); List<CABankAccount> closedBankAccounts = new ArrayList<>(); for (BankAccount account : items) { // business logic calling repository to few detail } super.write(activeBankAccounts); closedAccountsWriter.write(closedBankAccounts); } }
核心优化方案
1. 修复分区逻辑,实现真正的数据分片
当前分区仅分配了index,所有分区的Reader都读取全量数据,导致重复处理、资源竞争,这是性能瓶颈的核心。
- 优化分区器:根据数据分片键(如ID范围)将数据均匀拆分到各分区
protected final Partitioner partitioner() { return gridSize -> { Map<String, ExecutionContext> map = new HashMap<>(gridSize); // 先获取总数据量,计算分片范围 long totalCount = clientService.getTotalAccountCount(); long partitionRange = totalCount / gridSize; for (int i = 1; i <= gridSize; i++) { ExecutionContext context = new ExecutionContext(); long startId = (i - 1) * partitionRange + 1; // 最后一个分区处理剩余所有数据 long endId = i == gridSize ? totalCount : i * partitionRange; context.putLong("startId", startId); context.putLong("endId", endId); map.put("partition_" + i, context); } return map; }; }
- 修改Reader,根据分片的startId/endId查询对应子集,避免全量拉取:
// 在ClientBankAccountsItemReader中添加获取分片参数的逻辑 @Override protected synchronized void doOpen()throws Exception { if (!isOpen) { log.info("Opening partitioned Client Account Reader..."); SearchAccountRequest request = getSearchAccountRequest(); // 从StepExecutionContext获取分片范围 ExecutionContext stepContext = getExecutionContext(); request.setStartId(stepContext.getLong("startId")); request.setEndId(stepContext.getLong("endId")); clientBankAccounts = clientService.getAccountsApi(request).subscribe(this::addToQueue); this.isOpen = true; } }
2. 优化Reactive Reader的队列与阻塞逻辑
当前Reader存在队列无界、阻塞超时不合理、全量拉取的问题:
- 给队列设置合理容量(如1000),避免内存溢出:
private BlockingQueue<ClientBankAccount> queue = new LinkedBlockingQueue<>(1000);
- 调整
offer超时逻辑,改用put()确保数据不丢失:
private void addToQueue(ClientBankAccount ba) { try { queue.put(ba); } catch (InterruptedException e) { log.warn("INTERRUPTED! Could not ADD Client Bank Account to the internal queue!"); Thread.currentThread().interrupt(); } }
- 缩短
doRead()的阻塞时间,避免不必要的等待:
@Override protected ClientBankAccount doRead() { try { return this.queue.poll(1, SECONDS); } catch (InterruptedException e) { log.warn("INTERRUPTED! Could not READ client Bank Account from the internal queue!"); Thread.currentThread().interrupt(); } return null; }
3. 优化JpaItemWriter的批量处理逻辑
当前Writer存在循环单查、事务拆分的问题:
- 将循环单查替换为批量查询,减少DB交互次数:
@Override public void write(List<? extends BankAccount> items) { if (isEmpty(items)) { return; } List<BankAccount> activeBankAccounts = new ArrayList<>(); List<CABankAccount> closedBankAccounts = new ArrayList<>(); // 收集所有需要查询的关联ID Set<Long> relatedIds = items.stream() .map(BankAccount::getRelatedId) .collect(Collectors.toSet()); // 批量查询关联数据 Map<Long, RelatedEntity> relatedEntityMap = relatedRepository.findAllById(relatedIds).stream() .collect(Collectors.toMap(RelatedEntity::getId, entity -> entity)); // 业务逻辑使用批量查询结果 for (BankAccount account : items) { RelatedEntity entity = relatedEntityMap.get(account.getRelatedId()); // 处理active/closed账户分类逻辑 } // 批量写入 super.write(activeBankAccounts); closedAccountsWriter.write(closedBankAccounts); }
- 确保JpaItemWriter的
batchSize配置与chunk size匹配,开启批量插入/更新:
// 在Writer配置中设置batchSize @Bean public JpaItemWriter<BankAccount> bankAccountsItemWriter(EntityManagerFactory entityManagerFactory) { JpaItemWriter<BankAccount> writer = new JpaItemWriter<>(); writer.setEntityManagerFactory(entityManagerFactory); writer.setBatchSize(1000); // 与chunk size一致 return writer; }
4. 调整Chunk Size与线程池配置
- 将chunk size调整为1000,与batch-size匹配,减少事务提交次数
- 确保taskExecutor的核心线程数与grid size一致(如6),避免线程切换开销:
// 示例线程池配置 @Bean public TaskExecutor applicationTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(6); executor.setMaxPoolSize(6); executor.setQueueCapacity(100); executor.setThreadNamePrefix("batch-partition-"); executor.initialize(); return executor; }
- 移除或简化过于频繁的日志Listener(如
chunkLogger()),减少IO开销
5. H2数据库优化
- 使用内存模式:
jdbc:h2:mem:testdb;DB_CLOSE_DELAY=-1,避免磁盘IO - 关闭不必要的日志:添加
;LOG=0到JDBC URL - 配置锁模式:
;DEFAULT_LOCK_MODE=0,减少锁竞争
内容的提问来源于stack exchange,提问作者avinash
相关产品推荐
相关产品推荐

