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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 15:24:54