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

Spring Batch跨步骤数据传递问题排查:城市与历史数据匹配

Spring Batch跨步骤传递历史数据匹配城市失败排查

我开发了一个Spring Batch项目,分为两个步骤分别读取Cities和Histories两个TXT文件,ItemReader可正常读取数据,也能将两者合并写入单个JSON文件(这是预期目标)。但写入前需要在ItemProcessor中匹配城市与其对应的历史数据,让每个城市后跟随匹配的历史记录。尝试通过JobExecutionContext从ProcessStep1传递历史列表到ProcessStep2进行匹配,但未能成功实现,请求帮忙排查问题。


BatchConfiguration配置类

@Configuration
public class BatchConfiguration {

  @Bean
  DataSource dataSource() {
    return new EmbeddedDatabaseBuilder().setType(EmbeddedDatabaseType.H2).addScript("/org/springframework/batch/core/schema-h2.sql").generateUniqueName(
        true).build();
  }

  @Bean
  PlatformTransactionManager transactionManager(DataSource ds) {
    return new JdbcTransactionManager(ds);
  }

  @Bean
  public Job runjob(JobRepository jobRepository, PlatformTransactionManager txManager)
      throws Exception {
    return new JobBuilder("runJob", jobRepository)
        .start(step1(jobRepository, txManager))
        .next(step2(jobRepository, txManager))
        .build();
  }

  @Bean
  public Step step1(JobRepository jobRepository, PlatformTransactionManager txManager)
      throws Exception {
    return new StepBuilder("step1")
        .repository(jobRepository)
        .<Historique,Historique>chunk(100)
        .reader(readerStep1())
        .processor(processStep1())
        .writer(writerStep1())
        .listener(promotionListener())
        .transactionManager(txManager)
        .build();
  }

  @Bean
  public ExecutionContextPromotionListener promotionListener() {
    ExecutionContextPromotionListener listener = new ExecutionContextPromotionListener();
    listener.setKeys(new String[]{"histories"});
    return listener;
  }

  @Bean
  public Step step2(JobRepository jobRepository, PlatformTransactionManager txManager)
      throws Exception {
    return  new StepBuilder("step1")
        .repository(jobRepository)
        .<City,City>chunk(100)
        .reader(readerStep2())
        .processor(processStep2())
        .writer(writerStep2())
        .transactionManager(txManager)
        .build();
  }

  @Bean
  public ItemReader<Historique> readerStep1()
      throws Exception {
    FlatFileItemReader<Historique> readerStep1 = new FlatFileItemReader<Historique>();
    readerStep1.setLineMapper(new DefaultLineMapper() {{
      setLineTokenizer(new DelimitedLineTokenizer() {{
        setNames(new String[]{"mod", "date_eff", "typecom_av", "com_av", "tncc_av", "ncc_av", "nccenr_av", "libelle_av", "typecom_ap", "com_ap", "tncc_ap", "ncc_ap",
                             "nccenr_ap", "libelle_ap"});
      }});
      setFieldSetMapper(new BeanWrapperFieldSetMapper<Historique>() {{
        setTargetType(Historique.class);
      }});
    }});
    readerStep1.setResource(new ClassPathResource("documents/HistoriqueCities.txt"));
    return readerStep1;
  }

  @Bean
  public ItemReader<City> readerStep2()
      throws Exception {
    FlatFileItemReader<City> readerStep2 = new FlatFileItemReader<City>();
    readerStep2.setLineMapper(new DefaultLineMapper() {{
      setLineTokenizer(new DelimitedLineTokenizer() {{
        setNames(new String[]{"typecom", "com", "reg", "dep", "arr", "tncc", "ncc", "nccenr", "libelle", "can", "comparent"});
      }});
      setFieldSetMapper(new BeanWrapperFieldSetMapper<City>() {{
        setTargetType(City.class);
      }});
    }});
    readerStep2.setResource(new ClassPathResource("documents/Cities.txt"));
    return readerStep2;
  }

  @Bean
  public ProcessStep1 processStep1() {
    return new ProcessStep1();
  }

  @Bean
  public ProcessStep2 processStep2() {
    return new ProcessStep2();
  }

  @Bean
  public ItemWriter<Historique> writerStep1() {
    JsonFileItemWriter<Historique> writerStep1 = new JsonFileItemWriter<>(new FileSystemResource("src/main/java/output/testStep.json"),
        new JsonObjectMarshaller<Historique>() {
          @Override
          public String marshal(Historique object) {
            try {
              ObjectMapper mapper = new ObjectMapper();
              mapper.configure(SerializationFeature.INDENT_OUTPUT, true);
              return mapper.writeValueAsString(object);
            }
            catch(JsonProcessingException e) {
              e.printStackTrace();
              return null;
            }
          }
        }) {
      @Override
      public void open(ExecutionContext executionContext)
          throws ItemStreamException {
        super.open(executionContext);
        try {
          getOutputState().write("[{Historique: }");
        }
        catch(IOException e) {
          throw new RuntimeException(e);
        }
      }
    };
    writerStep1.setResource(new FileSystemResource("src/main/java/output/testStep.json"));
    writerStep1.setJsonObjectMarshaller(new JsonObjectMarshaller<Historique>() {

      @Override
      public String marshal(Historique historique) {
        Map<String, String> mapStory = new LinkedHashMap<>();
        mapStory.put("parent", historique.getCom_av());
        mapStory.put("typeMod", historique.getTypecom_av());
        mapStory.put("typeModNumber", historique.getMod());
        mapStory.put("oldLabel", historique.getLibelle_ap());
        mapStory.put("Label", historique.getLibelle_av());
        mapStory.put("date_eff", historique.getDate_eff());
        List<ProcessStep1> historiqueList = new ArrayList<>();


        return mapStory.toString();
      }
    });
    return writerStep1;

  }

  @Bean
  public ItemWriter<City> writerStep2() {
    JsonFileItemWriter<City> writerStep2 = new JsonFileItemWriter<>(new FileSystemResource("src/main/java/output/testStep.json"),
        new JsonObjectMarshaller<City>() {
          @Override
          public String marshal(City object) {
            try {
              ObjectMapper mapper = new ObjectMapper();
              mapper.enable(SerializationFeature.INDENT_OUTPUT);
              return mapper.writeValueAsString(object);
            }
            catch(JsonProcessingException e) {
              e.printStackTrace();
              return null;
            }
          }
        }) {

      @Override
      public void open(ExecutionContext executionContext)
          throws ItemStreamException {
        super.open(executionContext);
        try {
          getOutputState().write("[{City: }");
        }
        catch(IOException e) {
          throw new RuntimeException(e);
        }
      }
    };
    writerStep2.setResource(new FileSystemResource("src/main/java/output/testStep.json"));
    writerStep2.setAppendAllowed(true);
    writerStep2.setJsonObjectMarshaller(new JsonObjectMarshaller<City>() {
      @Override
      public String marshal(City city) {
        Map<String, String> mapCity = new LinkedHashMap<>();
        mapCity.put("parent", city.getCom());
        mapCity.put("typeMod", city.getTypecom());
        mapCity.put("typeModNumber", city.getLibelle());
        mapCity.put("oldLabel", city.getComparent());
        mapCity.put("Label", city.getCountry());
        System.out.println(mapCity.toString() + "test");
        return mapCity.toString();

      }
    });
    return writerStep2;
  }
}

ProcessStep1类

@Component
public class ProcessStep1 implements ItemProcessor<Historique, Historique>, StepExecutionListener {

  private DataHolder dataHolder;
  private ChunkContext chunkContext;

  @Override
  public Historique process(Historique historique) throws Exception {

    ExecutionContext executionContext = chunkContext.getStepContext().getStepExecution().getJobExecution().getExecutionContext();

    List<Historique> historiqueList = new ArrayList<>();
    executionContext.put("histories", historiqueList);
    return historique;
  }

  @Override
  public void beforeStep(StepExecution stepExecution) {
    StepExecutionListener.super.beforeStep(stepExecution);
    dataHolder = new DataHolder();
    stepExecution.getJobExecution().getExecutionContext().put("histories", dataHolder);

  }
  @Override
  public ExitStatus afterStep(StepExecution stepExecution) {
    DataHolder dataHolder = (DataHolder) stepExecution.getJobExecution().getExecutionContext().get("histories");
    return StepExecutionListener.super.afterStep(stepExecution);
  }
}

ProcessStep2类

@Component
public class ProcessStep2 implements ItemProcessor<City, City>, StepExecutionListener {

  private DataHolder dataHolder;
  private ChunkContext chunkContext;

  @Override
  public City process(City city) throws Exception {
    ExecutionContext executionContext = chunkContext.getStepContext().getStepExecution().getJobExecution().getExecutionContext();
    List<Historique> historiqueList = (List<Historique>) executionContext.get("histories");
    if(historiqueList == null) {
      historiqueList = new ArrayList<>();
    }
    city.setHistoriqueList(historiqueList);

return city;
  }

  @Override
  public void beforeStep(StepExecution stepExecution) {
    StepExecutionListener.super.beforeStep(stepExecution);
    DataHolder dataHolder = (DataHolder) stepExecution.getJobExecution().getExecutionContext().get("histories");
    stepExecution.getJobExecution().getExecutionContext().put("cities", dataHolder);
  }
  @Override
  public ExitStatus afterStep(StepExecution stepExecution) {
    return StepExecutionListener.super.afterStep(stepExecution);}

  }

DataHolder类

@Component
public class DataHolder implements Serializable {
  private Map<String, Object> dataMap;

  public DataHolder() {
    dataMap = new ConcurrentHashMap<>();
  }

  public Map<String, Object> getDataMap() {
    return dataMap;
  }

  public void setDataMap(Map<String, Object> dataMap) {
    this.dataMap = dataMap;
  }
}

问题排查与修复方案

1. 核心问题点

  • ChunkContext未注入:ProcessStep1和ProcessStep2中的chunkContext未通过Spring注入,直接调用会抛出空指针异常。
  • 历史数据存储逻辑错误:ProcessStep1的process方法每次新建空列表覆盖JobExecutionContext数据,导致最终无历史数据;且混合使用DataHolder和直接存List,逻辑混乱。
  • Step名称重复:BatchConfiguration中step2的StepBuilder使用了"step1"作为名称,会导致步骤定义冲突。
  • 类型转换错误:ProcessStep2中直接将DataHolder强转为List<Historique>,会抛出类型转换异常。
  • 数据收集与写入冲突:Step1的writer负责写入JSON,但我们需要先收集所有历史数据,写入操作会干扰数据传递逻辑。

2. 修复步骤

修复ProcessStep1:正确收集历史数据

@Component
public class ProcessStep1 implements ItemProcessor<Historique, Historique>, StepExecutionListener {

    private List<Historique> historiqueList;

    @Override
    public void beforeStep(StepExecution stepExecution) {
        // 初始化历史列表并存入JobExecutionContext
        historiqueList = new ArrayList<>();
        stepExecution.getJobExecution().getExecutionContext().put("histories", historiqueList);
    }

    @Override
    public Historique process(Historique historique) throws Exception {
        // 将每条历史数据添加到列表中
        historiqueList.add(historique);
        return historique;
    }

    @Override
    public ExitStatus afterStep(StepExecution stepExecution) {
        return StepExecutionListener.super.afterStep(stepExecution);
    }
}

修复ProcessStep2:正确获取并匹配历史数据

@Component
public class ProcessStep2 implements ItemProcessor<City, City>, StepExecutionListener {

    private List<Historique> historiqueList;

    @Override
    public void beforeStep(StepExecution stepExecution) {
        // 从JobContext中获取历史列表
        historiqueList = (List<Historique>) stepExecution.getJobExecution().getExecutionContext().get("histories");
        if (historiqueList == null) {
            historiqueList = new ArrayList<>();
        }
    }

    @Override
    public City process(City city) throws Exception {
        // 根据城市唯一标识(如com字段)匹配对应历史记录
        List<Historique> matchedHistories = historiqueList.stream()
                .filter(historique -> city.getCom().equals(historique.getCom_av()) || city.getCom().equals(historique.getCom_ap()))
                .collect(Collectors.toList());
        city.setHistoriqueList(matchedHistories);
        return city;
    }
}

修复BatchConfiguration中的Step名称

@Bean
public Step step2(JobRepository jobRepository, PlatformTransactionManager txManager) throws Exception {
    return new StepBuilder("step2") // 修改为唯一名称"step2"
            .repository(jobRepository)
            .<City,City>chunk(100)
            .reader(readerStep2())
            .processor(processStep2())
            .writer(writerStep2())
            .transactionManager(txManager)
            .build();
}

调整Step1的Writer(可选)

如果Step1仅用于收集数据,不需要写入JSON,可改为空实现:

@Bean
public ItemWriter<Historique> writerStep1() {
    return items -> {
        // 空实现,仅收集数据不写入文件
    };
}

确保序列化正确性

Historique和City类需实现Serializable接口,保证JobExecutionContext中的对象可序列化存储。

3. 验证逻辑

  • Step1执行时,ProcessStep1会收集所有历史数据到列表并存入JobExecutionContext。
  • Step2执行时,ProcessStep2在beforeStep阶段获取历史列表,为每个城市匹配对应历史记录。
  • 最终writerStep2将包含匹配历史的城市数据写入JSON文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 17:39:52