You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.08 16:57:35