Spring Batch 5中如何实现数据聚合并写入数据库?
Spring Batch 5 小时数据聚合为日值的实现方案
优先推荐:数据库层面聚合(适合百万级数据)
大数据量下,直接在SQL中完成聚合是性能最优的方案,避免应用层处理大量数据带来的内存和性能开销。
实现步骤
编写聚合SQL的ItemReader
用JdbcCursorItemReader直接读取聚合后的日值数据,返回单个DailyValue对象,完全匹配JdbcBatchItemWriter的处理逻辑:@Bean public JdbcCursorItemReader<DailyValue> dailyValueReader(DataSource dataSource) { // 按location_id和日期分组,SUM可替换为AVG/MAX等你需要的聚合逻辑 String aggregateSql = """ SELECT location_id, DATE(datetime) AS datetime, SUM(hourly_value) AS daily_value FROM hourly_value_table GROUP BY location_id, DATE(datetime) ORDER BY location_id, datetime """; return new JdbcCursorItemReaderBuilder<DailyValue>() .dataSource(dataSource) .sql(aggregateSql) .rowMapper((rs, rowNum) -> { DailyValue dailyValue = new DailyValue(); dailyValue.setLocationId(rs.getLong("location_id")); dailyValue.setDatetime(rs.getDate("datetime").toLocalDate()); dailyValue.setDailyValue(rs.getDouble("daily_value")); return dailyValue; }) .fetchSize(1000) // 批量读取,提升性能 .build(); }配置常规的JdbcBatchItemWriter
直接写入DailyValue对象即可,无需特殊处理:@Bean public JdbcBatchItemWriter<DailyValue> dailyValueWriter(DataSource dataSource) { String insertSql = """ INSERT INTO daily_value_table (location_id, datetime, daily_value) VALUES (:locationId, :datetime, :dailyValue) """; return new JdbcBatchItemWriterBuilder<DailyValue>() .dataSource(dataSource) .sql(insertSql) .beanMapped() .build(); }组装Job和Step
按常规Spring Batch配置组装,Chunk大小可根据数据库性能调整:@Bean public Job aggregateHourlyToDailyJob(JobRepository jobRepository, Step aggregateStep) { return new JobBuilder("aggregateHourlyToDailyJob", jobRepository) .start(aggregateStep) .build(); } @Bean public Step aggregateStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ItemReader<DailyValue> reader, ItemWriter<DailyValue> writer) { return new StepBuilder("aggregateStep", jobRepository) .<DailyValue, DailyValue>chunk(1000, transactionManager) .reader(reader) .writer(writer) .build(); }
备选方案:应用层聚合(适合复杂业务规则场景)
如果必须在应用层实现聚合逻辑(比如有SQL无法处理的复杂业务规则),需解决Chunk边界跨分组和类型匹配的问题:
核心思路
- 读取时必须按
location_id和datetime排序,确保同一分组的小时数据连续 - 自定义带状态的处理器,缓存当前分组的小时数据,当遇到新分组时输出聚合后的日值
- 实现
ItemStream接口,维护缓存状态,避免Job重启时数据不一致
实现示例
带状态的聚合处理器
@Component public class HourlyToDailyProcessor implements ItemProcessor<HourlyValue, DailyValue>, ItemStream { private Long currentLocationId; private LocalDate currentDate; private List<Double> hourlyValues = new ArrayList<>(); @Override public DailyValue process(HourlyValue item) throws Exception { LocalDate itemDate = item.getDatetime().toLocalDate(); // 首次处理或遇到新分组时,输出上一个分组的聚合结果 if (currentLocationId != null && (currentLocationId != item.getLocationId() || !currentDate.equals(itemDate))) { DailyValue aggregated = aggregate(currentLocationId, currentDate); // 更新当前分组信息 currentLocationId = item.getLocationId(); currentDate = itemDate; hourlyValues.clear(); hourlyValues.add(item.getHourlyValue()); return aggregated; } // 同一分组,累加数据 currentLocationId = item.getLocationId(); currentDate = itemDate; hourlyValues.add(item.getHourlyValue()); // 暂不输出,等待下一个分组 return null; } // 聚合逻辑,按需求修改(比如SUM/AVG等) private DailyValue aggregate(Long locationId, LocalDate date) { DailyValue dailyValue = new DailyValue(); dailyValue.setLocationId(locationId); dailyValue.setDatetime(date); dailyValue.setDailyValue(hourlyValues.stream().mapToDouble(Double::doubleValue).sum()); return dailyValue; } @Override public void open(ExecutionContext executionContext) throws ItemStreamException { // 从ExecutionContext恢复状态,支持Job重启 currentLocationId = executionContext.getLong("currentLocationId", null); currentDate = executionContext.getLocalDate("currentDate", null); hourlyValues = (List<Double>) executionContext.get("hourlyValues", new ArrayList<>()); } @Override public void update(ExecutionContext executionContext) throws ItemStreamException { // 保存状态到ExecutionContext executionContext.putLong("currentLocationId", currentLocationId); executionContext.putLocalDate("currentDate", currentDate); executionContext.put("hourlyValues", hourlyValues); } @Override public void close() throws ItemStreamException { // 处理最后一个分组的剩余数据,需配合监听器完成 } }添加StepExecutionListener处理最后一个分组
因为处理器在最后一个Item时不会输出结果,需要用监听器在Step结束时写入剩余的聚合数据:@Component public class AggregationCompletionListener implements StepExecutionListener { private final HourlyToDailyProcessor processor; private final JdbcBatchItemWriter<DailyValue> writer; public AggregationCompletionListener(HourlyToDailyProcessor processor, JdbcBatchItemWriter<DailyValue> writer) { this.processor = processor; this.writer = writer; } @Override public void afterStep(StepExecution stepExecution) throws Exception { // 处理最后一个分组的剩余数据 if (processor.currentLocationId != null) { DailyValue lastAggregate = processor.aggregate(processor.currentLocationId, processor.currentDate); writer.write(List.of(lastAggregate)); } } }配置Step时添加监听器
@Bean public Step aggregateStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ItemReader<HourlyValue> reader, HourlyToDailyProcessor processor, JdbcBatchItemWriter<DailyValue> writer, AggregationCompletionListener listener) { return new StepBuilder("aggregateStep", jobRepository) .<HourlyValue, DailyValue>chunk(1000, transactionManager) .reader(reader) .processor(processor) .writer(writer) .listener(listener) .build(); }
你之前的错误原因
你遇到的JdbcBatchItemWriter写入失败问题,本质是处理器输出类型与Writer输入类型不匹配:
JdbcBatchItemWriter<DailyValue>期望接收的是List<DailyValue>(即多个单个DailyValue对象)- 但你的处理器返回的是
List<DailyValue>作为单个Item,导致Writer试图把整个ArrayList当作一条数据库记录插入,自然报错。
正确的做法是让处理器输出单个DailyValue对象(或null,暂不输出),而非集合类型。
内容的提问来源于stack exchange,提问作者Niki Ryom Hansen
相关产品推荐
相关产品推荐

