Spring Batch如何循环执行步骤并传递前置步骤数据?
分层批处理:读取上游DB记录并逐个调用API的实现方案
针对你提出的批处理需求(大洲→国家→州的分层数据同步),步骤2和步骤3的核心逻辑是读取上一步存入DB的记录,提取ID作为参数调用下游API,再将结果写入DB。以下是具体实现思路和代码示例(以Spring Batch为例,这是Java生态最常用的批处理框架):
核心组件设计(以步骤2为例)
步骤2需要三个核心组件,职责明确:
- Reader:从DB读取步骤1存入的所有大洲记录,提取需要的字段(至少包含大洲ID)
- Processor:接收Reader传来的大洲对象,用其ID调用国家API,将返回的国家数据转换为DB实体
- Writer:将Processor输出的国家实体批量写入DB
1. 实现步骤2的Reader
使用数据库游标读取器(JdbcCursorItemReader)直接查询大洲表,获取所有需要的字段:
@Bean public JdbcCursorItemReader<Continent> continentReader(DataSource dataSource) { return new JdbcCursorItemReaderBuilder<Continent>() .dataSource(dataSource) .sql("SELECT id, name FROM continents") // 替换为你的大洲表名和字段 .rowMapper((rs, rowNum) -> { Continent continent = new Continent(); continent.setId(rs.getLong("id")); continent.setName(rs.getString("name")); return continent; }) .name("continentReader") .build(); }
如果大洲数据量极大,建议改用JdbcPagingItemReader实现分页读取,避免内存溢出。
2. 实现步骤2的Processor
在Processor中完成“用大洲ID调用API”和“数据转换”的逻辑,确保每个大洲记录触发一次API调用:
@Bean public ItemProcessor<Continent, List<Country>> countryApiProcessor(CountryApiClient countryApiClient) { // CountryApiClient是你封装的调用国家API的工具类 return continent -> { // 调用API获取该大洲下的所有国家 List<CountryApiResponse> apiResponses = countryApiClient.getCountriesByContinentId(continent.getId()); // 将API返回的DTO转换为DB实体 return apiResponses.stream() .map(response -> { Country country = new Country(); country.setContinentId(continent.getId()); country.setApiId(response.getCountryId()); country.setName(response.getCountryName()); // 补充其他字段映射 return country; }) .collect(Collectors.toList()); }; }
如果API返回单条数据而非列表,直接返回单个Country对象即可,无需转成列表。
3. 实现步骤2的Writer
用批量写入器将国家数据存入DB,提升写入效率:
@Bean public JdbcBatchItemWriter<Country> countryWriter(DataSource dataSource) { return new JdbcBatchItemWriterBuilder<Country>() .dataSource(dataSource) .sql("INSERT INTO countries (continent_id, api_id, name) VALUES (:continentId, :apiId, :name)") .beanMapped() // 自动映射实体字段到SQL参数 .build(); }
4. 组装步骤2的Step
将上述组件组装成Step,注意设置chunk(1)确保逐个处理大洲记录:
@Bean public Step step2(StepBuilderFactory stepBuilderFactory, ItemReader<Continent> continentReader, ItemProcessor<Continent, List<Country>> countryApiProcessor, ItemWriter<List<Country>> countryWriter) { return stepBuilderFactory.get("step2") .<Continent, List<Country>>chunk(1) // 每个chunk处理1个大洲,保证逐个调用API .reader(continentReader) .processor(countryApiProcessor) .writer(countryWriter) .build(); }
设置chunk(1)的原因是:每个大洲需要单独触发API调用,避免批量处理时API调用逻辑复杂化。如果API支持批量查询(比如一次性传入多个大洲ID),可以调大chunk值并修改Processor逻辑,提升效率。
步骤3的复用逻辑
步骤3(国家→州)的实现完全复用步骤2的思路:
- 用
JdbcCursorItemReader读取DB中的所有国家记录(提取国家ID) - 在Processor中用国家ID调用州API,转换为州实体
- 用批量写入器将州数据存入DB
- 组装Step时同样设置合适的chunk值
关键优化点
- 重试机制:API调用可能超时或失败,可通过Spring Retry给Processor添加重试逻辑,避免单次失败导致整个Step终止
- 事务管理:如果需要保证单个大洲的所有国家数据要么全成功要么全回滚,可在Step级别配置事务
- 错误处理:添加
ItemListener处理API调用失败的记录,比如写入错误日志表,后续单独重试
内容的提问来源于stack exchange,提问作者solver115
相关产品推荐
相关产品推荐

