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
相关产品推荐
相关产品推荐

