如何通过Spring Batch注解配置实现读取与批处理:两日文件价格差计算需求求助
Spring Batch实现两日产品价格差计算方案
你的问题核心是需要关联两个文件的同产品数据,但顺序执行的Step无法直接共享数据,这里给你一个清晰的解决方案:先把昨日的产品价格数据加载到内存Map中并存入Job上下文,然后在处理今日文件时,用这个Map来计算价格差。整个流程分为两个Step,完全符合Spring Batch的规范。
1. 定义实体类
首先创建两个实体类,分别对应文件中的产品价格数据,以及最终输出的价格差数据:
// 用于读取文件的产品价格实体 public class ProductPrice { private String productName; private int price; private String date; public ProductPrice() {} public ProductPrice(String productName, int price, String date) { this.productName = productName; this.price = price; this.date = date; } // 省略getter和setter方法,自行补充 } // 用于输出的价格差实体 public class PriceDifference { private String productName; private int difference; public PriceDifference() {} public PriceDifference(String productName, int difference) { this.productName = productName; this.difference = difference; } // 省略getter和setter方法,自行补充 }
2. 编写Tasklet加载昨日数据
用Tasklet来完成一次性的昨日文件读取操作,把产品名和价格的映射存入JobExecutionContext,供后续Step使用:
@Component public class YesterdayDataLoaderTasklet implements Tasklet { private final FlatFileItemReader<ProductPrice> yesterdayFileReader; public YesterdayDataLoaderTasklet(@Qualifier("yesterdayProductReader") FlatFileItemReader<ProductPrice> yesterdayFileReader) { this.yesterdayFileReader = yesterdayFileReader; } @Override public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { Map<String, Integer> yesterdayPriceMap = new HashMap<>(); ProductPrice productPrice; // 逐行读取昨日文件 while ((productPrice = yesterdayFileReader.read()) != null) { yesterdayPriceMap.put(productPrice.getProductName(), productPrice.getPrice()); } // 将映射存入Job上下文 chunkContext.getStepContext().getJobExecutionContext().put("yesterdayPriceMap", yesterdayPriceMap); return RepeatStatus.FINISHED; } }
3. 配置文件读取器
分别配置读取昨日和今日文件的FlatFileItemReader,处理CSV格式的文件:
@Configuration public class BatchReadersConfig { // 昨日文件读取器 @Bean public FlatFileItemReader<ProductPrice> yesterdayProductReader() { return new FlatFileItemReaderBuilder<ProductPrice>() .name("yesterdayProductReader") .resource(new FileSystemResource("your/path/to/yesterday_file.csv")) // 替换为实际路径 .linesToSkip(1) // 跳过表头行 .delimited() .delimiter(",") .names("productName", "price", "date") .fieldSetMapper(fieldSet -> new ProductPrice( fieldSet.readString("productName"), fieldSet.readInt("price"), fieldSet.readString("date") )) .build(); } // 今日文件读取器 @Bean public FlatFileItemReader<ProductPrice> todayProductReader() { return new FlatFileItemReaderBuilder<ProductPrice>() .name("todayProductReader") .resource(new FileSystemResource("your/path/to/today_file.csv")) // 替换为实际路径 .linesToSkip(1) // 跳过表头行 .delimited() .delimiter(",") .names("productName", "price", "date") .fieldSetMapper(fieldSet -> new ProductPrice( fieldSet.readString("productName"), fieldSet.readInt("price"), fieldSet.readString("date") )) .build(); } }
4. 编写价格差处理器
在Processor中从Job上下文取出昨日价格Map,计算今日与昨日的价格差:
@Component public class PriceDifferenceProcessor implements ItemProcessor<ProductPrice, PriceDifference> { @Override public PriceDifference process(ProductPrice todayProduct) throws Exception { // 获取Job上下文里的昨日价格映射 StepExecution stepExecution = StepSynchronizationManager.getContext().getStepExecution(); Map<String, Integer> yesterdayPriceMap = (Map<String, Integer>) stepExecution.getJobExecution().getExecutionContext().get("yesterdayPriceMap"); Integer yesterdayPrice = yesterdayPriceMap.get(todayProduct.getProductName()); if (yesterdayPrice == null) { // 如果昨日没有该产品,可以返回null让Spring Batch跳过这条记录,或者抛出异常 // return null; throw new IllegalArgumentException("未找到产品" + todayProduct.getProductName() + "的昨日价格"); } int difference = todayProduct.getPrice() - yesterdayPrice; return new PriceDifference(todayProduct.getProductName(), difference); } }
5. 编写结果文件写入器
配置FlatFileItemWriter来输出最终的价格差CSV文件:
@Component public class PriceDifferenceWriter extends FlatFileItemWriter<PriceDifference> { public PriceDifferenceWriter() { setResource(new FileSystemResource("your/path/to/output_difference.csv")); // 替换为输出路径 setHeaderCallback(writer -> writer.write("Productname,price")); // 写入表头 setLineAggregator(new DelimitedLineAggregator<>() {{ setDelimiter(","); setFieldExtractor(new BeanWrapperFieldExtractor<>() {{ setNames(new String[]{"productName", "difference"}); }}); }}); } }
6. 配置Job和Step
最后把所有组件组装成完整的Batch Job:
@Configuration @EnableBatchProcessing public class BatchJobConfig { private final JobBuilderFactory jobBuilderFactory; private final StepBuilderFactory stepBuilderFactory; private final YesterdayDataLoaderTasklet yesterdayDataLoaderTasklet; private final FlatFileItemReader<ProductPrice> todayProductReader; private final PriceDifferenceProcessor priceDifferenceProcessor; private final PriceDifferenceWriter priceDifferenceWriter; public BatchJobConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory, YesterdayDataLoaderTasklet yesterdayDataLoaderTasklet, @Qualifier("todayProductReader") FlatFileItemReader<ProductPrice> todayProductReader, PriceDifferenceProcessor priceDifferenceProcessor, PriceDifferenceWriter priceDifferenceWriter) { this.jobBuilderFactory = jobBuilderFactory; this.stepBuilderFactory = stepBuilderFactory; this.yesterdayDataLoaderTasklet = yesterdayDataLoaderTasklet; this.todayProductReader = todayProductReader; this.priceDifferenceProcessor = priceDifferenceProcessor; this.priceDifferenceWriter = priceDifferenceWriter; } // 第一步:加载昨日数据到Job上下文 @Bean public Step loadYesterdayDataStep() { return stepBuilderFactory.get("loadYesterdayDataStep") .tasklet(yesterdayDataLoaderTasklet) .build(); } // 第二步:处理今日数据,计算价格差并输出 @Bean public Step calculatePriceDifferenceStep() { return stepBuilderFactory.get("calculatePriceDifferenceStep") .<ProductPrice, PriceDifference>chunk(100) // 每100条数据作为一个块处理,可根据需求调整 .reader(todayProductReader) .processor(priceDifferenceProcessor) .writer(priceDifferenceWriter) .build(); } // 主Job @Bean public Job priceDifferenceJob() { return jobBuilderFactory.get("priceDifferenceJob") .start(loadYesterdayDataStep()) .next(calculatePriceDifferenceStep()) .build(); } }
关键注意事项
- 文件路径替换:请把代码中所有的文件路径替换为你实际的文件路径,也可以通过
@Value注解从配置文件中读取路径,更灵活。 - 异常处理:如果存在昨日没有的新产品,Processor中可以选择跳过(返回null)或者抛出异常,根据你的业务需求调整。
- 性能优化:对于每日1万条记录的场景,内存存储Map完全足够;如果是超大数据量,可以考虑把昨日数据存入数据库,处理今日数据时关联查询。
内容的提问来源于stack exchange,提问作者Codinghubby
相关产品推荐
相关产品推荐

