Spring Batch中如何审计CompositeItemWriter各子Writer的记录写入量?
嗨 Jeff,我之前在Spring Batch项目里也碰到过一模一样的问题——默认的StepListener只能追踪顶层Writer的指标,拿不到各个子Writer的写入量。下面给你几个可行的方案,既能追踪所有子Writer的计数,还能把审计信息存入DATA_AUDIT表:
方案一:给每个子Writer套一层计数包装器
这是最直观低耦合的方案,我们给CompositeItemWriter里的每个子Writer都加一个带计数功能的包装类,不修改原有Writer的业务逻辑,只负责记录写入数量。
步骤1:实现计数包装器
创建一个通用的CountingItemWriterWrapper,实现ItemWriter接口,内部持有真实的子Writer,用线程安全的原子变量记录写入数:
public class CountingItemWriterWrapper<T> implements ItemWriter<T> { private final ItemWriter<T> delegate; private final AtomicInteger writeCount = new AtomicInteger(0); public CountingItemWriterWrapper(ItemWriter<T> delegate) { this.delegate = delegate; } @Override public void write(List<? extends T> items) throws Exception { // 先执行原有写入逻辑 delegate.write(items); // 累加本次写入的记录数 writeCount.addAndGet(items.size()); } // 获取累计写入数 public int getWriteCount() { return writeCount.get(); } // 重置计数,避免下一次Step执行污染数据 public void resetCount() { writeCount.set(0); } }
步骤2:配置CompositeItemWriter
在配置CompositeWriter时,把每个子Writer都用包装器包裹,同时把这些包装器注入到自定义的审计监听器中:
@Bean public CompositeItemWriter<YourRecordType> compositeItemWriter() { CompositeItemWriter<YourRecordType> compositeWriter = new CompositeItemWriter<>(); // 给每个子Writer套上计数包装器 CountingItemWriterWrapper<YourRecordType> activeAccountsWriter = new CountingItemWriterWrapper<>(activeAccountsItemWriter()); CountingItemWriterWrapper<YourRecordType> inactiveAccountsWriter = new CountingItemWriterWrapper<>(inactiveAccountsItemWriter()); CountingItemWriterWrapper<YourRecordType> activeFlagWriter = new CountingItemWriterWrapper<>(activeFlagItemWriter()); // 其他子Writer同理 compositeWriter.setDelegates(List.of(activeAccountsWriter, inactiveAccountsWriter, activeFlagWriter)); return compositeWriter; }
步骤3:自定义审计监听器
创建一个StepExecutionListener,在Step执行完成后(afterStep方法),获取每个包装器的计数,写入DATA_AUDIT表:
@Component public class DataAuditStepListener implements StepExecutionListener { private final CountingItemWriterWrapper<YourRecordType> activeAccountsWriter; private final CountingItemWriterWrapper<YourRecordType> inactiveAccountsWriter; private final CountingItemWriterWrapper<YourRecordType> activeFlagWriter; private final JdbcTemplate jdbcTemplate; // 构造函数注入所有依赖 public DataAuditStepListener(CountingItemWriterWrapper<YourRecordType> activeAccountsWriter, CountingItemWriterWrapper<YourRecordType> inactiveAccountsWriter, CountingItemWriterWrapper<YourRecordType> activeFlagWriter, JdbcTemplate jdbcTemplate) { this.activeAccountsWriter = activeAccountsWriter; this.inactiveAccountsWriter = inactiveAccountsWriter; this.activeFlagWriter = activeFlagWriter; this.jdbcTemplate = jdbcTemplate; } @Override public ExitStatus afterStep(StepExecution stepExecution) { // 获取各子Writer的累计写入数 int activeCount = activeAccountsWriter.getWriteCount(); int inactiveCount = inactiveAccountsWriter.getWriteCount(); int flagCount = activeFlagWriter.getWriteCount(); int totalWriteCount = activeCount + inactiveCount + flagCount; // 写入DATA_AUDIT表 jdbcTemplate.update( "INSERT INTO DATA_AUDIT (step_name, read_count, total_write_count, active_accounts_count, inactive_accounts_count, active_flag_count, job_execution_id) " + "VALUES (?, ?, ?, ?, ?, ?, ?)", stepExecution.getStepName(), stepExecution.getReadCount(), totalWriteCount, activeCount, inactiveCount, flagCount, stepExecution.getJobExecutionId() ); // 重置计数,为下一次执行做准备 activeAccountsWriter.resetCount(); inactiveAccountsWriter.resetCount(); activeFlagWriter.resetCount(); return stepExecution.getExitStatus(); } }
方案二:利用ExecutionContext传递计数
如果不想额外包装Writer,可以让每个子Writer自己把写入数存入Step的ExecutionContext,之后在监听器里汇总这些数据。
步骤1:修改子Writer逻辑
给每个子Writer添加@BeforeStep注解注入StepExecution,在写入完成后更新ExecutionContext中的计数:
@Component public class ActiveAccountsItemWriter implements ItemWriter<YourRecordType> { private StepExecution stepExecution; // 注入StepExecution @BeforeStep public void setStepExecution(StepExecution stepExecution) { this.stepExecution = stepExecution; } @Override public void write(List<? extends YourRecordType> items) throws Exception { // 原有写入业务逻辑 // ... // 更新ExecutionContext中的计数 ExecutionContext context = stepExecution.getExecutionContext(); int currentCount = context.getInt("activeAccountsCount", 0); context.putInt("activeAccountsCount", currentCount + items.size()); } }
其他子Writer(比如InactiveAccountsItemWriter)同理,用不同的key(比如inactiveAccountsCount)记录自己的计数。
步骤2:在监听器中汇总计数
在自定义的StepExecutionListener的afterStep方法中,从ExecutionContext取出所有子Writer的计数,写入DATA_AUDIT表:
@Override public ExitStatus afterStep(StepExecution stepExecution) { ExecutionContext context = stepExecution.getExecutionContext(); int activeCount = context.getInt("activeAccountsCount", 0); int inactiveCount = context.getInt("inactiveAccountsCount", 0); int flagCount = context.getInt("activeFlagCount", 0); int totalWriteCount = activeCount + inactiveCount + flagCount; // 写入DATA_AUDIT表的逻辑同方案一 // ... // 注意:ExecutionContext会被持久化,所以如果需要重置,可以在afterStep里清空对应key context.remove("activeAccountsCount"); context.remove("inactiveAccountsCount"); context.remove("activeFlagCount"); return stepExecution.getExitStatus(); }
方案三:自定义带计数的CompositeItemWriter
如果想把计数逻辑统一封装,可以继承原有的CompositeItemWriter,扩展出计数功能:
public class CountingCompositeItemWriter<T> extends CompositeItemWriter<T> { // 用Map存储每个子Writer的计数,key可以是Writer的名称或标识 private final Map<String, AtomicInteger> writerCountMap = new HashMap<>(); // 存储子Writer到名称的映射,配置时设置 private final Map<ItemWriter<? super T>, String> writerNameMap; public CountingCompositeItemWriter(Map<ItemWriter<? super T>, String> writerNameMap) { this.writerNameMap = writerNameMap; } @Override public void write(List<? extends T> items) throws Exception { for (ItemWriter<? super T> delegate : getDelegates()) { String writerName = writerNameMap.get(delegate); delegate.write(items); // 累加计数 writerCountMap.computeIfAbsent(writerName, k -> new AtomicInteger(0)).addAndGet(items.size()); } } // 根据Writer名称获取计数 public int getCountByWriterName(String writerName) { return writerCountMap.getOrDefault(writerName, new AtomicInteger(0)).get(); } // 重置所有计数 public void resetAllCounts() { writerCountMap.values().forEach(count -> count.set(0)); } }
配置时,给每个子Writer关联一个名称,之后在监听器中注入这个自定义的CompositeWriter,获取各个子Writer的计数即可。
注意事项
- 如果使用多线程Step(比如配置了
TaskExecutor),一定要用线程安全的计数变量(比如AtomicInteger),或者用ThreadLocal记录每个线程的计数,最后汇总所有线程的结果。 - 每次Step执行完成后必须重置计数,避免下一次执行时带上之前的累计值。
- 审计记录要关联
job_execution_id,方便后续关联批处理作业的执行信息。
内容的提问来源于stack exchange,提问作者Jeff Cook

