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

