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

Jakarta EE Batch并行场景下如何批量处理ItemProcessor多数据?

Jakarta EE Batch 百万数据并行处理流程优化建议

问题背景

我正尝试基于Jakarta EE Batch实现一个数据处理流程:从数据库导入约100万条数据,后续分多步批量补充Web服务数据,第一步每批处理1万条、第二步每批15万条,每步均用5线程并行执行。

单步配置及单条数据处理无问题,但尚未找到在并行(partition=5)的ItemProcessor中可靠获取并处理多条数据的方法。

流程预期为:

  1. 从DB获取数据;
  2. 启动5次步骤1,每批处理1万条直至所有数据处理完成;
  3. 启动5次步骤2,每批处理15万条直至所有数据处理完成;
  4. 导出数据。

当前配置示例

<step id="getWsDataA" next="getWsDataB">
    <chunk>
        <reader ref="getWsDataAItemReader"/>
        <processor ref="getWsDataAItemMockProcessor"/>
        <writer ref="getWsDataAItemJpaWriter"/>
    </chunk>
    <partition>
        <plan partitions="5"></plan>
    </partition>
</step>

<step id="getWsDataB" next="export">
    <chunk>
        <reader ref="getWsDataAItemReader"/>
        <processor ref="getWsDataAItemMockProcessor"/>
        <writer ref="getWsDataAItemJpaWriter"/>
    </chunk>
    <partition>
        <plan partitions="5"></plan>
    </partition>
</step>

<step id="export">
    <batchlet ref="ExportBatchlet">
        <properties>
            <property name="path" value="C:\Temp\"/>
        </properties>
    </batchlet>
</step>

现有实现代码

Job配置

<?xml version="1.0" encoding="UTF-8"?>
<job id="hugeImport" xmlns="https://jakarta.ee/xml/ns/jakartaee" version="2.0">
    <step id="dummyItems" next="chunkProcessor">
        <batchlet ref="dummyItemsBatchlet">
            <properties>
                <property name="numberOfDummyItems" value="10"/>
            </properties>
        </batchlet>
    </step>

    <step id="chunkProcessor" next="reloadItemsQueue_001">
        <chunk>
            <reader ref="itemReader">
                <properties>
                    <property name="numberOfItems" value="2"/>
                </properties>
            </reader>
            <processor ref="itemMockProcessor"/>
            <writer ref="itemJpaWriter"/>
        </chunk>
        <partition>
            <plan partitions="2"></plan>
        </partition>
    </step>

    <step id="reloadItemsQueue_001" next="chunkProcessortest">
        <batchlet ref="reloadItemQueueBatchlet">
        </batchlet>
    </step>

    <step id="chunkProcessortest">
        <chunk>
            <reader ref="itemReader">
                <properties>
                    <property name="numberOfItems" value="3"/>
                </properties>
            </reader>
            <processor ref="itemMockProcessor"/>
            <writer ref="itemJpaWriter"/>
        </chunk>
        <partition>
            <plan partitions="2"></plan>
        </partition>
    </step>
</job>

数据实体类

public class ImportItem {
    private Long id;
    private String name;

    public String getName() {
        return name;
    }

    public void setName(String name) {
        this.name = name;
    }

    public ImportItem(long id, String name) {
        this.id = id;
        this.name = name;
    }

    @Override
    public String toString() {
        return "ImportItem{" + "id=" + id + ", name=" + name + '}';
    }
}

上下文管理类

import jakarta.batch.runtime.context.JobContext;
import jakarta.inject.Inject;
import jakarta.inject.Named;

import java.util.*;
import java.util.concurrent.ConcurrentLinkedQueue;

@Named
public class ImportJobContext {
    @Inject
    private JobContext jobContext;

    private final Queue<ImportItem> itemsToDo = new ConcurrentLinkedQueue<>();
    private final Queue<ImportItem> itemsForNextStep = new ConcurrentLinkedQueue<>();

    public void addItems(List<ImportItem> items) {
        getImportJobContext().itemsToDo.addAll(items);
    }

    public synchronized void reloadQueue(){
        getImportJobContext().itemsToDo.clear();
        getImportJobContext().itemsToDo.addAll(getImportJobContext().itemsForNextStep);
        getImportJobContext().itemsForNextStep.clear();
    }

    public synchronized List<ImportItem> getItems(int count) {
        List<ImportItem> items = new ArrayList<>(count);
        for (int i = 0; i < count; i++) {
            var item = getImportJobContext().itemsToDo.poll();

            if(item == null) {
                continue;
            }

            items.add(item);
            getImportJobContext().itemsForNextStep.add(item);
        }

        return items.isEmpty() ? null : items;
    }

    private ImportJobContext getImportJobContext() {
        if (jobContext.getTransientUserData() == null) {
            jobContext.setTransientUserData(this);
        }
        return (ImportJobContext) jobContext.getTransientUserData();
    }
}

Batchlet实现

import jakarta.batch.api.AbstractBatchlet;
import jakarta.batch.api.BatchProperty;
import jakarta.batch.runtime.BatchStatus;
import jakarta.inject.Inject;
import jakarta.inject.Named;

import java.util.ArrayList;
import java.util.List;

@Named
public class DummyItemsBatchlet extends AbstractBatchlet {

    @Inject
    private ImportJobContext jobContext;

    @Inject
    @BatchProperty
    private String numberOfDummyItems;

    @Override
    public String process() throws Exception {

        List<ImportItem> list = new ArrayList<>();
        for(int i=0; i<Integer.parseInt(numberOfDummyItems); i++){
            list.add(new ImportItem(i, "dummyItem" + i));
        }
        jobContext.addItems(list);

        return BatchStatus.COMPLETED.name();
    }
}
@Named
public class ReloadItemQueueBatchlet  extends AbstractBatchlet {

    @Inject
    private ImportJobContext jobContext;

    @Override
    public String process() throws Exception {
        System.out.println("ReloadItemQueueBatchlet.process");
        jobContext.reloadQueue();

        return BatchStatus.COMPLETED.name();
    }
}

Chunk组件实现

import jakarta.batch.api.BatchProperty;
import jakarta.batch.api.chunk.AbstractItemReader;
import jakarta.inject.Inject;
import jakarta.inject.Named;

import java.util.List;

@Named
public class ItemReader  extends AbstractItemReader {

    @Inject
    ImportJobContext importJobContext;

    @Inject
    @BatchProperty
    private String numberOfItems;

    @Override
    public List<ImportItem> readItem() throws Exception {

        int numberOfWorkerItems = 2;
        if(numberOfItems != null){
            numberOfWorkerItems = Integer.parseInt(numberOfItems);
        }

        return importJobContext.getItems(numberOfWorkerItems);
    }
}
import jakarta.batch.api.chunk.ItemProcessor;
import jakarta.inject.Named;

@Named
public class ItemMockProcessor implements ItemProcessor {

    @Override
    public Object processItem(Object o) throws Exception {
        System.out.println("--> processing " + o);
        return o;
    }
}
import jakarta.batch.api.chunk.AbstractItemWriter;
import jakarta.inject.Named;

import java.util.List;

@Named
public class ItemJpaWriter  extends AbstractItemWriter {

    @Override
    public void writeItems(List<Object> list) throws Exception {

        for (Object obj : list) {
            List<ImportItem> item = (List<ImportItem>) obj;
            System.out.println("--> Persisting " + item);
        }
    }
}

优化建议

1. 贴合Chunk模型设计,正确处理批量数据

当前ItemReader返回List<ImportItem>,导致ItemProcessor每次处理整个列表,违背了Chunk模型"读-处理-写"的原子性设计。调整方案:

  • 让ItemReader返回单条ImportItem,通过配置Chunk的batch-size控制每批处理数量(第一步设为10000,第二步设为150000)
  • 若需批量调用Web服务减少请求次数,可在ItemProcessor中累积数据,达到阈值后调用服务,或在ItemWriter阶段统一批量处理(更符合Chunk模型)

2. 优化并行分区的数据分片策略

当前全局队列抢数据的方式易导致负载不均,线程安全控制复杂。改为:

  • 基于数据库主键或范围分片,每个分区处理固定范围的数据(比如按ID分段:分区1处理1-20万,分区2处理20万-40万...)
  • 通过PartitionMapper动态计算每个分区的参数(起始ID、结束ID),传递给对应Reader
  • 示例分区配置调整:
<step id="getWsDataA" next="getWsDataB">
    <chunk batch-size="10000">
        <reader ref="getWsDataAItemReader">
            <properties>
                <property name="startId" value="#{partitionPlan['startId']}"/>
                <property name="endId" value="#{partitionPlan['endId']}"/>
            </properties>
        </reader>
        <processor ref="getWsDataAItemProcessor"/>
        <writer ref="getWsDataAItemJpaWriter"/>
    </chunk>
    <partition>
        <mapper ref="dataPartitionMapper"/>
    </partition>
</step>

3. 移除冗余的队列中转逻辑

内存缓存百万级数据易引发OOM,建议:

  • 每步处理后直接持久化到数据库(或临时表),下一个步骤的Reader直接从数据库读取已处理数据
  • 若必须内存传递,仅用JobContext的persistentUserData存储数据标识(如已处理ID列表),而非整个对象

4. 线程安全与并发优化

  • 移除ImportJobContext中不必要的synchronized方法,ConcurrentLinkedQueue本身是线程安全的,无需额外同步
  • 调整getItems方法:
public List<ImportItem> getItems(int count) {
    List<ImportItem> items = new ArrayList<>(count);
    for (int i = 0; i < count; i++) {
        ImportItem item = itemsToDo.poll();
        if (item == null) break;
        items.add(item);
        itemsForNextStep.add(item);
    }
    return items.isEmpty() ? null : items;
}

5. 代码结构与可读性优化

  • 为每个步骤创建独立的Reader/Processor/Writer,避免复用导致的配置混乱
  • 给ImportItem添加无参构造函数,方便框架实例化
  • 移除未使用的ImportItems类,简化代码结构

6. 错误处理与监控

  • 在Chunk配置中添加skip-limit和skip-policy,处理Web服务调用失败等异常
  • 利用JobOperator监控步骤执行状态、处理计数,添加日志记录每批数据处理情况

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 16:07:33