JBeret实现JSR-352中JobContext transientUserData跨步骤传递失败问题
问题分析与解决方案
你说得太对了,JobContext的命名确实很容易误导人——它根本不是全局作业级的上下文,而是每个步骤实例(包括每个分区的独立步骤实例)的本地上下文。每个分区都会创建自己的JobContext副本,所以它们的transientUserData完全隔离,而且当step1的所有分区执行完成后,这些本地的临时数据不会自动传递到step2的JobContext里,这就是你在step2拿到null的核心原因。
下面是JSR-352规范中,跨步骤、跨分区共享作业数据的标准方式,按推荐度排序:
1. 使用JobExecutionContext的executionContext(最推荐)
JobExecutionContext才是真正的全局作业上下文,它的executionContext是一个支持持久化的键值对存储(即使作业重启也能保留数据),对所有步骤、所有分区实例可见。
修改你的ItemWriter代码:
@Named public class EccWriter extends AbstractItemWriter { @Inject Logger logger; @Inject JobContext jobContext; @Override public void writeItems(List<Object> list) throws Exception { // 获取作业级的全局ExecutionContext ExecutionContext jobExecContext = jobContext.getJobExecution().getExecutionContext(); // 从全局上下文获取已处理列表,不存在则初始化 @SuppressWarnings("unchecked") ArrayList<String> processed = Optional.ofNullable( jobExecContext.get("processedUsers") ).map(ArrayList.class::cast).orElse(new ArrayList<>()); list.stream().map(UserLogin.class::cast).forEach(input -> { if (someConditions) { processed.add(input.getUserId()); } }); // 将更新后的列表放回全局上下文 jobExecContext.put("processedUsers", processed); } }
在step2的Batchlet中获取数据:
@Named public class EccMailBatchlet extends AbstractBatchlet { @Inject JobContext jobContext; @Override public String process() throws Exception { ExecutionContext jobExecContext = jobContext.getJobExecution().getExecutionContext(); ArrayList<String> processedUsers = (ArrayList<String>) jobExecContext.get("processedUsers"); // 这里就能拿到所有分区处理后的合并列表了 logger.info("Processed users: " + processedUsers); return "COMPLETED"; } }
注意事项:
ExecutionContext中的数据会被持久化(比如存储到数据库),所以存储的对象必须实现Serializable接口,否则会抛出序列化异常。- 多个分区同时修改同一数据时,要注意线程安全——可以用同步块,或者直接使用线程安全的集合(比如
CopyOnWriteArrayList)。
2. 使用CDI的@ApplicationScoped bean(适合单JVM作业)
如果你的作业运行在同一个JVM进程中,可以用CDI全局作用域bean来共享内存级数据。这种方式不需要持久化,适合存储临时中间结果。
创建全局共享bean:
@ApplicationScoped public class JobSharedData { // 用同步集合保证线程安全 private final List<String> processedUsers = Collections.synchronizedList(new ArrayList<>()); public void addProcessedUser(String userId) { processedUsers.add(userId); } // 返回副本避免外部直接修改内部集合 public List<String> getProcessedUsers() { return new ArrayList<>(processedUsers); } // 作业结束后清理数据,避免影响下一次执行 public void clear() { processedUsers.clear(); } }
在ItemWriter中注入使用:
@Named public class EccWriter extends AbstractItemWriter { @Inject Logger logger; @Inject JobSharedData sharedData; @Override public void writeItems(List<Object> list) throws Exception { list.stream().map(UserLogin.class::cast).forEach(input -> { if (someConditions) { sharedData.addProcessedUser(input.getUserId()); } }); } }
在step2的Batchlet中获取:
@Named public class EccMailBatchlet extends AbstractBatchlet { @Inject JobSharedData sharedData; @Override public String process() throws Exception { List<String> processedUsers = sharedData.getProcessedUsers(); logger.info("Processed users: " + processedUsers); // 作业完成后清理数据 sharedData.clear(); return "COMPLETED"; } }
注意事项:
- 这种方式只适合同一JVM环境,分布式作业(多节点运行)下每个节点的
@ApplicationScopedbean是独立的,数据不会共享。
3. 分区场景进阶:使用PartitionReducer合并结果
如果你的分区需要更规范的结果合并,可以使用PartitionReducer——它会在所有分区执行完成后统一处理所有分区的结果,再合并到作业级上下文。
实现PartitionReducer:
@Named public class EccPartitionReducer implements PartitionReducer { @Inject JobContext jobContext; @Override public void reduce(Collection<StepExecution> stepExecutions) throws Exception { ExecutionContext jobExecContext = jobContext.getJobExecution().getExecutionContext(); ArrayList<String> allProcessed = new ArrayList<>(); // 遍历所有分区的StepExecution,收集各自的处理结果 for (StepExecution stepExecution : stepExecutions) { ExecutionContext stepExecContext = stepExecution.getExecutionContext(); ArrayList<String> partitionProcessed = (ArrayList<String>) stepExecContext.get("partitionProcessed"); if (partitionProcessed != null) { allProcessed.addAll(partitionProcessed); } } jobExecContext.put("processedUsers", allProcessed); } }
在作业配置中指定Reducer:
<step id="step1" next="step2"> <chunk item-count="#{jobParameters['chunksize']}?:3"> <reader ref="eccReader"> </reader> <writer ref="eccWriter" /> </chunk> <partition> <mapper ref="eccMapper"> <properties> <property name="threads" value="#{jobParameters['threads']}?:3"/> <property name="records" value="#{jobParameters['records']}?:30"/> </properties> </mapper> <reducer ref="eccPartitionReducer"/> <!-- 添加这一行 --> </partition> </step>
这种方式更符合JSR-352的分区设计规范,每个分区将结果存到自己的StepExecution上下文,再由Reducer统一合并到作业级上下文。
内容的提问来源于stack exchange,提问作者Fabrizio Stellato
相关产品推荐
相关产品推荐

