Spring Batch分区与并行处理实现多门店订单数据同步方案咨询
用Spring Batch分区+并行步骤实现门店订单数据同步流程
整体流程设计
整个流程拆分为两个独立的分区Step,按顺序执行:
- 分区并行调用第一个REST API,获取5个门店的orderId总数及orderIds,保存到仓库
- 分区并行读取仓库中对应门店的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
相关产品推荐
相关产品推荐

