Spring Batch有状态Processor如何聚合多行数据生成目标业务对象
Spring Batch多行数据聚合实现方案
完全可以实现多行聚合生成目标对象的需求,以下是两种落地性最高的实现方式:
方案1:自定义聚合ItemReader(推荐)
该方案将聚合逻辑封装在Reader层,符合Spring Batch组件单一职责的设计规范,后续Processor、Writer直接接收聚合完成的Person对象即可,无需额外改造。
实现前提
你的查询SQL必须按照分组字段(firstName、lastName)排序,保证同一个人的所有薪资行连续返回,否则聚合逻辑会失效。
代码实现
- 首先定义数据库行对应的临时映射对象:
public class PersonSalaryRow { private String firstName; private String lastName; private String week; private Integer salary; // 省略getter、setter方法 }
- 自定义聚合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
相关产品推荐
相关产品推荐

