Spring Batch:如何从ItemListenerSupport将错误数据存入ExecutionContext
如何将ItemListenerSupport中累积的错误信息存入Spring Batch的Job ExecutionContext?
问题描述
我正在扩展ItemListenerSupport以捕获Spring Batch读取、处理、写入步骤中遇到的错误,当前的onWriteError代码如下:
@Override public void onWriteError(Exception ex, List<? extends BaseDomainDataObject> items) { logger.error("Encountered error on write", ex); String msgBody = ExceptionUtils.getStackTrace(ex); numProcessedMap.computeIfAbsent("numErrors", val -> items.size()); errorMap.put(numErrors.addAndGet(1), msgBody); }
请问如何将map中累积的所有错误信息存入Step或Job(优先选择Job)的ExecutionContext中?
解决方案
好问题!要把你累积的错误信息存入Job ExecutionContext(这个选择很合理,因为Job级别的上下文可以跨多个Step共享,后续在Job结束后也能方便获取),你需要借助Spring Batch的StepExecutionListener接口来关联到Job上下文,具体实现步骤如下:
让你的监听器实现
StepExecutionListener接口
这个接口能让你在Step执行前后获取到StepExecution实例,而通过StepExecution可以直接拿到关联的JobExecution,进而访问Job级别的ExecutionContext。保存
StepExecution实例供后续使用
在beforeStep方法中把当前的StepExecution保存为类成员变量,这样在onWriteError或者afterStep里就能随时调用。将累积的错误数据存入Job ExecutionContext
推荐在afterStep方法中统一存入数据(批量操作更高效),当然你也可以在每次onWriteError触发时实时写入,根据你的需求选择即可。
下面是完整的代码示例:
import org.springframework.batch.core.ExitStatus; import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.StepExecutionListener; import org.springframework.batch.core.listener.ItemListenerSupport; import org.apache.commons.lang3.exception.ExceptionUtils; import java.util.List; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class CustomErrorTrackingListener extends ItemListenerSupport<BaseDomainDataObject, BaseDomainDataObject> implements StepExecutionListener { private static final Logger logger = LoggerFactory.getLogger(CustomErrorTrackingListener.class); private StepExecution stepExecution; // 用线程安全的集合处理多线程场景下的错误累积 private final ConcurrentHashMap<String, Integer> numProcessedMap = new ConcurrentHashMap<>(); private final ConcurrentHashMap<Integer, String> errorMap = new ConcurrentHashMap<>(); private final AtomicInteger numErrors = new AtomicInteger(0); @Override public void beforeStep(StepExecution stepExecution) { // 保存当前Step的执行实例 this.stepExecution = stepExecution; } @Override public ExitStatus afterStep(StepExecution stepExecution) { // 获取Job级别的ExecutionContext ExecutionContext jobExecutionContext = stepExecution.getJobExecution().getExecutionContext(); // 将累积的错误数据存入Job上下文 jobExecutionContext.put("batchErrorMap", errorMap); jobExecutionContext.put("batchErrorStats", numProcessedMap); jobExecutionContext.put("totalErrorCount", numErrors.get()); return ExitStatus.COMPLETED; } @Override public void onWriteError(Exception ex, List<? extends BaseDomainDataObject> items) { logger.error("Encountered error during write operation", ex); String errorStackTrace = ExceptionUtils.getStackTrace(ex); // 累积错误统计和详情 numProcessedMap.compute("numErrors", (key, currentValue) -> currentValue == null ? items.size() : currentValue + items.size()); errorMap.put(numErrors.addAndGet(1), errorStackTrace); } }
关键细节说明
- 线程安全:使用
ConcurrentHashMap和AtomicInteger确保在多线程Step(比如并行Chunk处理)场景下,错误累积操作不会出现并发问题。 - Job vs Step ExecutionContext:如果后续需要在其他Step或者Job结束后查看错误信息,Job级别的上下文是最佳选择;如果只需要当前Step内访问,也可以改用
stepExecution.getExecutionContext()存入Step上下文。 - 数据读取:后续要获取这些错误数据时,只需要从
JobExecution.getExecutionContext()中取出对应key的值即可,比如在JobListener或者其他组件中访问。
内容的提问来源于stack exchange,提问作者aatuc210
相关产品推荐
相关产品推荐

