Spring Batch中如何根据处理器结果运行时动态选择ItemWriter
Spring Batch 按ItemProcessor处理结果动态切换ItemWriter实现方式
Spring Batch默认Chunk执行逻辑中,ItemProcessor抛出未捕获异常会直接触发当前Chunk回滚、Step终止,不会进入ItemWriter执行链路,要实现校验通过走成功写入器、校验失败走异常写入器的需求,有两种可落地的实现方案:
方案一:封装带状态的处理结果 + 分类复合写入器(推荐,事务一致性最优)
该方案不依赖容错跳过机制,所有数据都走完整的Chunk处理链路,事务边界统一,适合生产环境使用:
- 第一步:定义统一的处理结果封装类,标记数据校验状态,同时携带原始数据、异常信息
// 校验处理结果封装 public class ProcessResult { private boolean valid; private Sample data; private ValidationException validationError; // 省略全参构造、getter/setter public static ProcessResult success(Sample data) { return new ProcessResult(true, data, null); } public static ProcessResult fail(Sample data, ValidationException e) { return new ProcessResult(false, data, e); } }
- 第二步:调整原有Bean校验处理器,内部捕获校验异常不向外抛出,返回对应状态的ProcessResult对象
public class BeanValidationProcessor implements ItemProcessor<Sample, ProcessResult> { @Override public ProcessResult process(Sample item) { try { // 原有Bean校验逻辑 validate(item); return ProcessResult.success(item); } catch (ValidationException e) { return ProcessResult.fail(item, e); } } }
- 第三步:分别实现两个子写入器
- 成功写入器:筛选
valid=true的结果,提取原始Sample数据执行正常持久化逻辑 - 异常写入器:筛选
valid=false的结果,提取原始数据、校验异常信息执行异常记录逻辑
- 成功写入器:筛选
- 第四步:配置Spring Batch自带的
ClassifierCompositeItemWriter作为主写入器,按结果状态路由到对应子写入器
@Bean public ItemWriter<ProcessResult> routeWriter( ItemWriter<ProcessResult> successWriter, ItemWriter<ProcessResult> errorWriter ) { ClassifierCompositeItemWriter<ProcessResult> compositeWriter = new ClassifierCompositeItemWriter<>(); // 配置路由规则 compositeWriter.setClassifier(result -> result.isValid() ? successWriter : errorWriter); return compositeWriter; }
- 第五步:调整Step配置,匹配链路泛型,注册子写入器到Step生命周期
@Bean public Step step( ItemReader<Sample> itemReader, BeanValidationProcessor beanProcessor, ItemWriter<ProcessResult> routeWriter, ItemWriter<ProcessResult> successWriter, ItemWriter<ProcessResult> errorWriter ) { return stepBuilderFactory.get("our step") .<Sample, ProcessResult>chunk(1) .reader(itemReader) .processor(beanProcessor) .writer(routeWriter) // 注册子写入器,保证生命周期回调、事务正常生效 .stream(successWriter) .stream(errorWriter) .build(); }
方案二:保留Processor抛异常逻辑 + 容错跳过机制+监听器分流
如果不想修改原有Processor抛出校验异常的逻辑,可以通过Spring Batch的容错跳过能力兜住异常,配合监听器实现分流:
@Bean public Step step( ItemReader<Sample> itemReader, ItemProcessor<Sample, Sample> beanProcessor, ItemWriter<Sample> successWriter, ItemWriter<Sample> errorWriter ) { return stepBuilderFactory.get("our step") .<Sample, Sample>chunk(1) .reader(itemReader) .processor(beanProcessor) .writer(successWriter) // 开启容错,配置校验异常为可跳过异常,不触发Chunk回滚 .faultTolerant() .skip(ValidationException.class) .skipLimit(Integer.MAX_VALUE) // 注册跳过监听器,捕获Processor抛出的校验异常,写入异常数据 .listener(new SkipListener<Sample, Sample>() { @Override public void onSkipInRead(Throwable t) {} @Override public void onSkipInWrite(Sample item, Throwable t) {} @Override public void onSkipInProcess(Sample item, Throwable t) { try { errorWriter.write(List.of(item)); } catch (Exception e) { throw new RuntimeException("异常数据写入失败", e); } } }) .stream(errorWriter) .build(); }
注意事项:
- 你当前配置的chunk大小为1,两种方案都不会出现单chunk内部分数据成功/失败的事务冲突;如果后续调大chunk size,优先选择第一种方案,事务边界更统一,不会出现异常数据写入和正常数据写入事务不一致的问题。
- 使用复合写入器时,必须通过
.stream()方法将所有子写入器注册到Step中,否则子写入器的生命周期回调(open、update、close)不会被触发,文件、数据库连接类的写入器会出现运行异常。
内容的提问来源于stack exchange,提问作者Vinicius de Andrade
相关产品推荐
相关产品推荐

