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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 02:47:31