Jakarta EE Batch并行场景下如何批量处理ItemProcessor多数据?
Jakarta EE Batch 百万数据并行处理流程优化建议
问题背景
我正尝试基于Jakarta EE Batch实现一个数据处理流程:从数据库导入约100万条数据,后续分多步批量补充Web服务数据,第一步每批处理1万条、第二步每批15万条,每步均用5线程并行执行。
单步配置及单条数据处理无问题,但尚未找到在并行(partition=5)的ItemProcessor中可靠获取并处理多条数据的方法。
流程预期为:
- 从DB获取数据;
- 启动5次步骤1,每批处理1万条直至所有数据处理完成;
- 启动5次步骤2,每批处理15万条直至所有数据处理完成;
- 导出数据。
当前配置示例
<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
相关产品推荐
相关产品推荐

