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

Spring Batch按院系输出学生记录及批次汇总实现问询

Spring Batch 按院系分组处理并输出汇总信息的实现方案

核心思路

由于你的数据已按院系ID排序,核心方案是在处理过程中跟踪当前院系的状态:每当遇到新院系ID时,先输出上一个院系的汇总信息,再处理新院系的学生记录;Step结束时处理最后一个院系的汇总。无需复杂监听器组合,用「带状态的ItemProcessor + @AfterStep注解」即可实现需求。

具体实现步骤

1. 定义数据模型

明确两个核心实体:

  • Student:存储学生基础信息(院系ID、学生ID、姓名、缴费金额等)
  • DepartmentSummary:存储院系汇总数据(院系ID、学生总数、总缴费金额)
public class Student {
    private String deptId;
    private String studentId;
    private String name;
    private BigDecimal feeAmount;
    // getter、setter省略
}

public class DepartmentSummary {
    private String deptId;
    private int studentCount;
    private BigDecimal totalFee;
    // 构造器、getter、setter省略
}

2. 实现带状态的ItemProcessor

用@StepScope保证处理器在Step内的状态隔离(避免多Step共享状态),在处理器中维护当前院系的统计数据,遇到新院系时触发汇总写入:

@StepScope
@Component
public class StudentProcessingProcessor implements ItemProcessor<Student, Object> {

    private String currentDeptId;
    private int studentCount;
    private BigDecimal totalFee;
    private final ItemWriter<Object> delegateWriter;

    // 注入通用写入器,用于输出汇总信息
    public StudentProcessingProcessor(ItemWriter<Object> delegateWriter) {
        this.delegateWriter = delegateWriter;
    }

    @Override
    public Object process(Student student) throws Exception {
        // 首次处理或遇到新院系
        if (currentDeptId == null || !currentDeptId.equals(student.getDeptId())) {
            // 非首次处理时,先写入上一个院系的汇总
            if (currentDeptId != null) {
                DepartmentSummary summary = new DepartmentSummary(currentDeptId, studentCount, totalFee);
                delegateWriter.write(Collections.singletonList(summary));
            }
            // 重置当前院系状态
            currentDeptId = student.getDeptId();
            studentCount = 0;
            totalFee = BigDecimal.ZERO;
        }
        // 更新当前院系统计数据
        studentCount++;
        totalFee = totalFee.add(student.getFeeAmount());
        // 返回学生记录,交给后续写入器处理
        return student;
    }

    // Step结束时处理最后一个院系的汇总
    @AfterStep
    public ExitStatus afterStep(StepExecution stepExecution) throws Exception {
        if (currentDeptId != null) {
            DepartmentSummary summary = new DepartmentSummary(currentDeptId, studentCount, totalFee);
            delegateWriter.write(Collections.singletonList(summary));
        }
        return stepExecution.getExitStatus();
    }
}

3. 配置多类型兼容的ItemWriter

因为需要同时写入Student和DepartmentSummary,用ClassifierCompositeItemWriter根据类型路由到不同的写入逻辑,实现同一文件内先写学生记录、再写汇总信息:

@Configuration
public class BatchWriterConfig {

    // 通用分类写入器
    @Bean
    public ClassifierCompositeItemWriter<Object> studentAndSummaryWriter() {
        ClassifierCompositeItemWriter<Object> writer = new ClassifierCompositeItemWriter<>();
        writer.setClassifier(item -> {
            if (item instanceof Student) {
                return studentItemWriter();
            } else if (item instanceof DepartmentSummary) {
                return summaryItemWriter();
            }
            throw new IllegalArgumentException("未知数据类型:" + item.getClass());
        });
        return writer;
    }

    // 学生记录写入器(CSV格式)
    @Bean
    public FlatFileItemWriter<Student> studentItemWriter() {
        FlatFileItemWriter<Student> writer = new FlatFileItemWriter<>();
        writer.setResource(new FileSystemResource("students_output.csv"));
        // 配置CSV行聚合规则
        writer.setLineAggregator(new DelimitedLineAggregator<>() {{
            setDelimiter(",");
            setFieldExtractor(new BeanWrapperFieldExtractor<>() {{
                setNames(new String[]{"deptId", "studentId", "name", "feeAmount"});
            }});
        }});
        // 写入前清空文件(可选,根据需求调整)
        writer.setShouldDeleteIfExists(true);
        return writer;
    }

    // 汇总信息写入器(追加到同一文件,加前缀区分)
    @Bean
    public FlatFileItemWriter<DepartmentSummary> summaryItemWriter() {
        FlatFileItemWriter<DepartmentSummary> writer = new FlatFileItemWriter<>();
        writer.setResource(new FileSystemResource("students_output.csv"));
        writer.setAppendAllowed(true); // 追加模式
        // 汇总行加前缀,方便识别
        writer.setLineSeparator(System.lineSeparator() + "汇总,");
        writer.setLineAggregator(new DelimitedLineAggregator<>() {{
            setDelimiter(",");
            setFieldExtractor(new BeanWrapperFieldExtractor<>() {{
                setNames(new String[]{"deptId", "studentCount", "totalFee"});
            }});
        }});
        return writer;
    }
}

4. 配置Repository Reader和Step

确保Repository Reader的查询结果按院系ID排序(以JpaPagingItemReader为例):

@Bean
public JpaPagingItemReader<Student> studentRepositoryReader(EntityManagerFactory emf) {
    JpaPagingItemReader<Student> reader = new JpaPagingItemReader<>();
    reader.setEntityManagerFactory(emf);
    reader.setQueryString("SELECT s FROM Student s ORDER BY s.deptId");
    reader.setPageSize(100); // 分页大小根据数据量调整
    reader.setSort(Collections.singletonMap("deptId", Sort.Direction.ASC));
    return reader;
}

// 配置Step
@Bean
public Step studentProcessingStep(StepBuilderFactory stepBuilderFactory,
                                  JpaPagingItemReader<Student> studentReader,
                                  StudentProcessingProcessor processor,
                                  ClassifierCompositeItemWriter<Object> writer) {
    return stepBuilderFactory.get("studentProcessingStep")
            .<Student, Object>chunk(100)
            .reader(studentReader)
            .processor(processor)
            .writer(writer)
            .build();
}

关于监听器的说明

这里用@AfterStep注解实现了StepExecutionListener的功能,用于处理最后一个院系的汇总。如果不需要在处理器中维护状态,也可以用ChunkListener结合全局状态,但处理器结合@AfterStep的方式更直接,且@StepScope保证了每个Step实例的处理器状态独立,避免并发问题。

内容的提问来源于stack exchange,提问作者Ashish Gupta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 00:35:45