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

Spring Batch分区与并行处理实现多门店订单数据同步方案咨询

用Spring Batch分区+并行步骤实现门店订单数据同步流程

整体流程设计

整个流程拆分为两个独立的分区Step,按顺序执行:

  1. 分区并行调用第一个REST API,获取5个门店的orderId总数及orderIds,保存到仓库
  2. 分区并行读取仓库中对应门店的orderId数据,调用第二个REST API获取订单列表,保存到仓库

核心组件实现细节

1. 第一个Step:获取并保存OrderId数据

分区器(Partitioner)

实现一个Partitioner,将5个门店作为独立分区,每个分区携带对应门店ID:

public class StorePartitioner implements Partitioner {
    @Override
    public Map<String, ExecutionContext> partition(int gridSize) {
        Map<String, ExecutionContext> partitions = new HashMap<>();
        List<String> storeIds = Arrays.asList("store1", "store2", "store3", "store4", "store5");
        for (String storeId : storeIds) {
            ExecutionContext context = new ExecutionContext();
            context.putString("storeId", storeId);
            partitions.put("partition-" + storeId, context);
        }
        return partitions;
    }
}

工作Step(Worker Step)

每个分区执行的具体逻辑:调用API、解析数据、写入仓库:

@Bean
public Step fetchOrderIdsStep(StepBuilderFactory stepBuilderFactory,
                              ItemReader<OrderIdData> orderIdReader,
                              ItemWriter<OrderIdData> orderIdWriter) {
    return stepBuilderFactory.get("fetchOrderIdsStep")
            .<OrderIdData, OrderIdData>chunk(1)
            .reader(orderIdReader)
            .writer(orderIdWriter)
            .build();
}

// 从REST API读取对应门店的OrderId数据
public class OrderIdRestReader implements ItemReader<OrderIdData> {
    private RestTemplate restTemplate;
    private String storeId;
    private boolean readCompleted = false;

    @Override
    public OrderIdData read() throws Exception {
        if (readCompleted) return null;
        ResponseEntity<OrderIdData> response = restTemplate.getForEntity(
                "http://api.example.com/stores/{storeId}/order-ids",
                OrderIdData.class,
                storeId
        );
        readCompleted = true;
        return response.getBody();
    }

    @BeforeStep
    public void beforeStep(StepExecution stepExecution) {
        this.storeId = stepExecution.getExecutionContext().getString("storeId");
    }
}

// 将OrderId数据写入仓库
public class OrderIdRepositoryWriter implements ItemWriter<OrderIdData> {
    private OrderIdRepository orderIdRepository;

    @Override
    public void write(List<? extends OrderIdData> items) throws Exception {
        orderIdRepository.save(items.get(0));
    }
}

主Step(Master Step)

配置分区并行执行工作Step:

@Bean
public Step masterFetchOrderIdsStep(StepBuilderFactory stepBuilderFactory,
                                    StorePartitioner storePartitioner,
                                    Step fetchOrderIdsStep,
                                    TaskExecutor taskExecutor) {
    return stepBuilderFactory.get("masterFetchOrderIdsStep")
            .partitioner(fetchOrderIdsStep.getName(), storePartitioner)
            .step(fetchOrderIdsStep)
            .taskExecutor(taskExecutor)
            .gridSize(5)
            .build();
}

2. 第二个Step:获取并保存订单列表数据

分区器(复用StorePartitioner)

和第一个Step的分区器逻辑一致,按门店ID划分分区。

工作Step(Worker Step)

每个分区执行的逻辑:读取仓库中的OrderId数据、调用API获取订单列表、写入仓库:

@Bean
public Step fetchOrderListStep(StepBuilderFactory stepBuilderFactory,
                               ItemReader<List<Order>> orderListReader,
                               ItemWriter<List<Order>> orderListWriter) {
    return stepBuilderFactory.get("fetchOrderListStep")
            .<List<Order>, List<Order>>chunk(1)
            .reader(orderListReader)
            .writer(orderListWriter)
            .build();
}

// 先读仓库数据,再调用API获取订单列表
public class OrderListRestReader implements ItemReader<List<Order>> {
    private RestTemplate restTemplate;
    private OrderIdRepository orderIdRepository;
    private String storeId;
    private boolean readCompleted = false;

    @Override
    public List<Order> read() throws Exception {
        if (readCompleted) return null;
        OrderIdData orderIdData = orderIdRepository.findByStoreId(storeId);
        if (orderIdData == null) {
            readCompleted = true;
            return Collections.emptyList();
        }
        ResponseEntity<List<Order>> response = restTemplate.postForEntity(
                "http://api.example.com/stores/{storeId}/orders",
                orderIdData,
                new ParameterizedTypeReference<List<Order>>() {},
                storeId
        );
        readCompleted = true;
        return response.getBody();
    }

    @BeforeStep
    public void beforeStep(StepExecution stepExecution) {
        this.storeId = stepExecution.getExecutionContext().getString("storeId");
    }
}

// 将订单列表批量写入仓库
public class OrderRepositoryWriter implements ItemWriter<List<Order>> {
    private OrderRepository orderRepository;

    @Override
    public void write(List<? extends List<Order>> items) throws Exception {
        orderRepository.saveAll(items.get(0));
    }
}

主Step(Master Step)

配置分区并行:

@Bean
public Step masterFetchOrderListStep(StepBuilderFactory stepBuilderFactory,
                                     StorePartitioner storePartitioner,
                                     Step fetchOrderListStep,
                                     TaskExecutor taskExecutor) {
    return stepBuilderFactory.get("masterFetchOrderListStep")
            .partitioner(fetchOrderListStep.getName(), storePartitioner)
            .step(fetchOrderListStep)
            .taskExecutor(taskExecutor)
            .gridSize(5)
            .build();
}

3. 作业(Job)组装

将两个主Step按顺序组装成完整Job:

@Bean
public Job storeOrderSyncJob(JobBuilderFactory jobBuilderFactory,
                             Step masterFetchOrderIdsStep,
                             Step masterFetchOrderListStep) {
    return jobBuilderFactory.get("storeOrderSyncJob")
            .start(masterFetchOrderIdsStep)
            .next(masterFetchOrderListStep)
            .build();
}

关键注意点

  • TaskExecutor配置:推荐用ThreadPoolTaskExecutor,核心线程数设为门店数(5),避免API调用限流
  • 异常处理:给工作Step添加监听器,处理API调用失败、仓库写入失败的场景,支持重试或标记异常门店
  • 数据隔离:仓库中存储的OrderId数据必须按门店ID区分,避免跨门店数据混淆
  • 分区参数传递:通过ExecutionContext传递门店ID,确保每个Worker Step只处理对应门店的数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 23:30:28