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

Spring Batch 5中如何实现数据聚合并写入数据库?

Spring Batch 5 小时数据聚合为日值的实现方案

优先推荐:数据库层面聚合(适合百万级数据)

大数据量下,直接在SQL中完成聚合是性能最优的方案,避免应用层处理大量数据带来的内存和性能开销。

实现步骤

  1. 编写聚合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();
    }
    
  2. 配置常规的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();
    }
    
  3. 组装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重启时数据不一致

实现示例

  1. 带状态的聚合处理器

    @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 {
            // 处理最后一个分组的剩余数据,需配合监听器完成
        }
    }
    
  2. 添加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));
            }
        }
    }
    
  3. 配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 20:57:13