Spring Batch跨步骤复用已读取ProductId的实现方案问询
问题
我开发了一个Spring Batch产品处理更新任务,用于为打印产品标签补充产品数据,核心步骤为productCategorizationStep和productDetailsConfigStep。规则是无详情或分类信息的产品会被过滤,但需要将所有已处理(包括未完成更新)的ProductId传递到最终的productProcessingSummaryStep生成汇总CSV文件,且不想重复使用ItemReader。已知可以通过StepExecution传递步骤间数据,但需要实现跨多步骤传递至最终步骤的方案。
附上当前作业配置代码:
import com.healthpartners.carelocation.batch.listener.JobResultListener; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.batch.core.*; import org.springframework.batch.core.configuration.annotation.JobBuilderFactory; import org.springframework.batch.core.launch.JobLauncher; import org.springframework.batch.core.launch.support.RunIdIncrementer; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.annotation.Scheduled; import java.util.Date; import java.util.UUID; @Slf4j @Configuration @RequiredArgsConstructor public class JobConfiguration { @Configuration private final JobBuilderFactory jobBuilderFactory; private final JobLauncher productAppJobLauncher; private final static String PRODUCT_JOB = "ProductJob"; private final JobResultListener jobResultListener; @Qualifier("productUpdateStep") private final Step productUpdateStep; @Qualifier("productDetailsConfigStep") private final Step productDetailsConfigStep; @Qualifier("productCategorizationStep") private final Step productCategorizationStep; /** * This step collects all the Products for which product details were added * by productId. It also displays in a file the products for which details could not * be found like below * PRODUCT_ID DETAILS_UPDATED CATEGORIZED * 12344 true true * 342222 false true * 32122 true false */ @Qualifier("productProcessingSummaryStep") private final Step productProcessingSummaryStep; @Bean("productProcessingJob") public Job processingJob() { return jobBuilderFactory.get(CARE_LOCATION_JOB) .incrementer(new RunIdIncrementer()) .listener(jobResultListener) .start(productUpdateStep) .next(productDetailsConfigStep) .next(productCategorizationStep) .next(productProcessingSummaryStep) .build(); } @Scheduled(cron = "${processing.schedule}") //3AM? public void runJob() throws Exception { JobParameters jobParameters = new JobParametersBuilder() .addString("jobId", UUID.randomUUID().toString()) .addDate("date", new Date()) .addLong("time", System.currentTimeMillis()).toJobParameters(); JobExecution execution = productAppJobLauncher.run(processingJob(), jobParameters); log.info("STATUS :: " + execution.getStatus()); } }
解决方案
要跨多步骤共享所有已处理的Product数据,核心是使用JobExecutionContext(作业级别的上下文,整个作业生命周期内共享),而非仅StepExecutionContext(步骤级上下文,仅当前步骤有效)。具体实现分四步:
1. 定义处理状态跟踪DTO
创建一个DTO来记录每个Product的处理结果,包含汇总CSV需要的字段:
import lombok.AllArgsConstructor; import lombok.Data; @Data @AllArgsConstructor public class ProductProcessingRecord { private String productId; private boolean detailsUpdated; private boolean categorized; }
2. 在处理步骤中收集数据到JobExecutionContext
在productDetailsConfigStep和productCategorizationStep中,通过自定义StepListener或ItemListener跟踪每个Product的处理状态,并将记录写入JobExecutionContext。
以productDetailsConfigStep为例,实现数据收集逻辑:
import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.StepSynchronizationManager; import org.springframework.batch.item.ItemProcessor; import java.util.ArrayList; import java.util.List; // 假设这是productDetailsConfigStep的配置类 public class ProductDetailsConfigStepConfig { // 从JobExecutionContext中获取或初始化记录列表 private List<ProductProcessingRecord> getProcessingRecords(StepExecution stepExecution) { List<ProductProcessingRecord> records = (List<ProductProcessingRecord>) stepExecution.getJobExecution() .getExecutionContext() .get("productProcessingRecords"); if (records == null) { records = new ArrayList<>(); stepExecution.getJobExecution().getExecutionContext().put("productProcessingRecords", records); } return records; } // 自定义ItemProcessor,同时更新处理状态 public ItemProcessor<Product, Product> detailsConfigProcessor() { return product -> { StepExecution stepExecution = StepSynchronizationManager.getContext().getStepExecution(); List<ProductProcessingRecord> records = getProcessingRecords(stepExecution); boolean detailsUpdated = false; // 执行详情补充逻辑,更新detailsUpdated状态 if (/* 详情补充成功判断逻辑 */) { detailsUpdated = true; } // 记录当前步骤的处理状态(后续分类步骤会更新categorized字段) records.add(new ProductProcessingRecord(product.getId(), detailsUpdated, false)); return detailsUpdated ? product : null; // 过滤未补充详情的产品 }; } // 步骤配置示例 public Step productDetailsConfigStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, ItemReader<Product> productReader, ItemWriter<Product> productWriter) { return new StepBuilder("productDetailsConfigStep", jobRepository) .<Product, Product>chunk(100, transactionManager) .reader(productReader) .processor(detailsConfigProcessor()) .writer(productWriter) .build(); } }
同理,在productCategorizationStep的ItemProcessor中,更新对应Product的分类状态:
import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.StepSynchronizationManager; import org.springframework.batch.item.ItemProcessor; import java.util.List; public ItemProcessor<Product, Product> categorizationProcessor() { return product -> { StepExecution stepExecution = StepSynchronizationManager.getContext().getStepExecution(); List<ProductProcessingRecord> records = (List<ProductProcessingRecord>) stepExecution.getJobExecution() .getExecutionContext() .get("productProcessingRecords"); boolean categorized = false; // 执行分类逻辑,更新categorized状态 if (/* 分类成功判断逻辑 */) { categorized = true; } // 找到对应ProductId的记录并更新分类状态 records.stream() .filter(record -> record.getProductId().equals(product.getId())) .findFirst() .ifPresent(record -> record.setCategorized(categorized)); return categorized ? product : null; // 过滤未分类的产品 }; }
3. 在汇总步骤中读取JobExecutionContext的数据
productProcessingSummaryStep不需要重复调用ItemReader,直接从JobExecutionContext读取收集到的所有ProductProcessingRecord,然后写入CSV:
import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.StepSynchronizationManager; import org.springframework.batch.item.ItemReader; import org.springframework.batch.item.file.FlatFileItemWriter; import org.springframework.batch.item.file.transform.BeanWrapperFieldExtractor; import org.springframework.batch.item.file.transform.DelimitedLineAggregator; import org.springframework.core.io.FileSystemResource; // 假设这是productProcessingSummaryStep的配置类 public class ProductProcessingSummaryStepConfig { // 自定义ItemReader,从JobExecutionContext读取数据 public ItemReader<ProductProcessingRecord> summaryReader() { return new ItemReader<>() { private List<ProductProcessingRecord> records; private int index = 0; @Override public ProductProcessingRecord read() throws Exception { if (records == null) { StepExecution stepExecution = StepSynchronizationManager.getContext().getStepExecution(); records = (List<ProductProcessingRecord>) stepExecution.getJobExecution() .getExecutionContext() .get("productProcessingRecords"); if (records == null) { records = new ArrayList<>(); } } return index < records.size() ? records.get(index++) : null; } }; } // CSV Writer配置 public FlatFileItemWriter<ProductProcessingRecord> summaryWriter() { FlatFileItemWriter<ProductProcessingRecord> writer = new FlatFileItemWriter<>(); writer.setResource(new FileSystemResource("product-processing-summary.csv")); writer.setHeaderCallback(writer1 -> writer1.write("PRODUCT_ID,DETAILS_UPDATED,CATEGORIZED")); DelimitedLineAggregator<ProductProcessingRecord> lineAggregator = new DelimitedLineAggregator<>(); lineAggregator.setDelimiter(","); BeanWrapperFieldExtractor<ProductProcessingRecord> fieldExtractor = new BeanWrapperFieldExtractor<>(); fieldExtractor.setNames(new String[]{"productId", "detailsUpdated", "categorized"}); lineAggregator.setFieldExtractor(fieldExtractor); writer.setLineAggregator(lineAggregator); return writer; } // 步骤配置 public Step productProcessingSummaryStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) { return new StepBuilder("productProcessingSummaryStep", jobRepository) .<ProductProcessingRecord, ProductProcessingRecord>chunk(100, transactionManager) .reader(summaryReader()) .writer(summaryWriter()) .build(); } }
4. 确保JobExecutionContext数据持久化
默认情况下,Spring Batch会将JobExecutionContext的数据持久化到作业元数据存储(如数据库),跨步骤共享无需额外配置。如果是本地内存作业,需使用MapJobRepositoryFactoryBean来保存上下文数据。
注意事项
- 若作业为分片(Partitioned)运行,需使用线程安全的集合(如
CopyOnWriteArrayList)存储记录,避免并发问题。 - 处理超大数量Product时,内存存储可能导致OOM,可改为临时文件存储中间数据,汇总时读取临时文件。
内容的提问来源于stack exchange,提问作者BreenDeen

