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

Spring Batch有状态Processor如何聚合多行数据生成目标业务对象

Spring Batch多行数据聚合实现方案

完全可以实现多行聚合生成目标对象的需求,以下是两种落地性最高的实现方式:


方案1:自定义聚合ItemReader(推荐)

该方案将聚合逻辑封装在Reader层,符合Spring Batch组件单一职责的设计规范,后续Processor、Writer直接接收聚合完成的Person对象即可,无需额外改造。

实现前提

你的查询SQL必须按照分组字段(firstName、lastName)排序,保证同一个人的所有薪资行连续返回,否则聚合逻辑会失效。

代码实现

  1. 首先定义数据库行对应的临时映射对象:
public class PersonSalaryRow {
    private String firstName;
    private String lastName;
    private String week;
    private Integer salary;
    // 省略getter、setter方法
}
  1. 自定义聚合Reader,封装你已经实现的JdbcCursorItemReader:
@StepScope
@Component
public class PersonAggregateReader implements ItemStreamReader<Person> {
    private final JdbcCursorItemReader<PersonSalaryRow> delegateReader;
    // 缓存上一轮读取到的跨行分组数据
    private PersonSalaryRow lastCachedRow;

    public PersonAggregateReader(JdbcCursorItemReader<PersonSalaryRow> delegateReader) {
        this.delegateReader = delegateReader;
    }

    @Override
    public Person read() throws Exception {
        Person currentPerson = null;
        // 优先读取上一轮缓存的跨行数据作为当前分组的起始
        if (lastCachedRow != null) {
            currentPerson = convertRowToPerson(lastCachedRow);
            lastCachedRow = null;
        }

        PersonSalaryRow currentRow;
        while ((currentRow = delegateReader.read()) != null) {
            // 首次读取初始化当前Person对象
            if (currentPerson == null) {
                currentPerson = convertRowToPerson(currentRow);
                continue;
            }
            // 分组匹配:同一个人则追加薪资数据
            if (isSamePersonGroup(currentPerson, currentRow)) {
                currentPerson.getPayment().add(convertRowToSalary(currentRow));
            } else {
                // 分组不匹配,缓存当前行,返回已聚合完成的Person
                lastCachedRow = currentRow;
                return currentPerson;
            }
        }
        // 所有行读取完成,返回最后一个聚合完成的Person
        return currentPerson;
    }

    // 判断是否属于同一个人的分组
    private boolean isSamePersonGroup(Person person, PersonSalaryRow row) {
        return person.getFirstName().equals(row.getFirstName())
                && person.getLastName().equals(row.getLastName());
    }

    // 从行数据初始化Person对象
    private Person convertRowToPerson(PersonSalaryRow row) {
        Person person = new Person();
        person.setFirstName(row.getFirstName());
        person.setLastName(row.getLastName());
        List<Salary> paymentList = new ArrayList<>();
        paymentList.add(convertRowToSalary(row));
        person.setPayment(paymentList);
        return person;
    }

    // 从行数据转换为Salary对象
    private Salary convertRowToSalary(PersonSalaryRow row) {
        Salary salary = new Salary();
        salary.setWeek(row.getWeek());
        salary.setSalary(row.getSalary());
        return salary;
    }

    // 委托原Reader的状态管理方法,支持失败重启
    @Override
    public void open(ExecutionContext executionContext) throws ItemStreamException {
        delegateReader.open(executionContext);
    }

    @Override
    public void update(ExecutionContext executionContext) throws ItemStreamException {
        delegateReader.update(executionContext);
    }

    @Override
    public void close() throws ItemStreamException {
        delegateReader.close();
    }
}

优势

  • 逻辑内聚,Reader直接输出完整业务对象,下游组件无感知
  • 天然支持chunk模式,无需额外处理边界数据
  • @StepScope保证步骤级缓存隔离,无线程安全问题

方案2:有状态ItemProcessor实现

如果不想修改已有的Reader逻辑,也可以在Processor层做聚合:

  • 给自定义Processor添加@StepScope注解,内部缓存当前正在聚合的Person对象
  • 每次接收单条行数据时,判断分组是否匹配:匹配则追加薪资,不匹配则返回上一个聚合完成的Person,更新缓存为当前行对应的新Person
  • 实现StepExecutionListener接口,在afterStep回调中将最后缓存的Person补充到输出结果中,避免最后一条数据丢失。
    该方案的缺点是需要额外处理边界数据,逻辑耦合度高于自定义Reader方案。

内容的提问来源于stack exchange,提问作者carlos palma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 06:54:03