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

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的思路:

  1. 用JdbcCursorItemReader读取DB中的所有国家记录(提取国家ID)
  2. 在Processor中用国家ID调用州API,转换为州实体
  3. 用批量写入器将州数据存入DB
  4. 组装Step时同样设置合适的chunk值

关键优化点

  • 重试机制:API调用可能超时或失败,可通过Spring Retry给Processor添加重试逻辑,避免单次失败导致整个Step终止
  • 事务管理:如果需要保证单个大洲的所有国家数据要么全成功要么全回滚,可在Step级别配置事务
  • 错误处理:添加ItemListener处理API调用失败的记录,比如写入错误日志表,后续单独重试

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 23:05:19