如何将Spring Batch分区数据按顺序传递给Tasklet?
Spring Batch 用Tasklet实现分区数据传递调用API
你现在的需求是用提前准备好的客户对象列表,拿每个对象的customerId调用API,已经写好接收customerId的Tasklet和ListPartitioner,但不知道怎么把分区后的数据传给Tasklet。先看你提供的chunk模式代码(翻译后):
@Bean(name="asyncStep") protected Step asyncStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) throws Exception { return new StepBuilder("myjob", jobRepository) .<EmployeeDTO,EmployeeDTO>chunk(2,transactionManager) // 按每2条数据为一个块处理 .reader(itemReader(null)) // 读取数据 .processor(asyncItemProcessor()) // 处理数据 .writer(asyncItemWriter()) // 写入数据 .build(); }
chunk模式是框架帮你封装了读写处理的流程,而Tasklet模式需要自己手动处理数据传递,下面是具体实现步骤:
把原始客户列表传入作业上下文
因为列表在作业启动前就准备好了,启动作业时可以把列表(或者如果列表存在缓存/数据库里,传个标识比如缓存key)放到JobParameters里,这样分区器和Tasklet都能拿到。让分区器把拆分后的数据存入分区上下文
修改你的ListPartitioner,把分区后的客户ID列表放到每个分区的ExecutionContext中,这样每个Tasklet实例都能拿到对应分区的数据:
public class CustomerListPartitioner implements Partitioner { @Override public Map<String, ExecutionContext> partition(int gridSize) { // 从JobParameters取出提前准备好的客户列表 List<Customer> customerList = (List<Customer>) JobParameterUtils.getJobParameter("customerList"); // 按指定大小拆分列表 List<List<Customer>> partitionedLists = splitList(customerList, gridSize); Map<String, ExecutionContext> partitions = new HashMap<>(); for (int i = 0; i < partitionedLists.size(); i++) { ExecutionContext context = new ExecutionContext(); // 把当前分区的customerId集合存入上下文 List<String> customerIds = partitionedLists.get(i).stream() .map(Customer::getCustomerId) .collect(Collectors.toList()); context.put("partitionedCustomerIds", customerIds); partitions.put("partition-" + i, context); } return partitions; } // 工具方法:拆分列表为多个子列表 private List<List<Customer>> splitList(List<Customer> list, int chunkSize) { List<List<Customer>> result = new ArrayList<>(); for (int i = 0; i < list.size(); i += chunkSize) { int endIndex = Math.min(i + chunkSize, list.size()); result.add(list.subList(i, endIndex)); } return result; } }
- 在Tasklet中读取分区数据并调用API
修改你的Tasklet,从当前Step的上下文里取出分区后的customerId列表,遍历调用API即可:
public class CustomerApiCallTasklet implements Tasklet { private final ApiClient apiClient; // 注入你用来调用API的客户端 @Override public RepeatStatus execute(StepContribution contribution, ChunkContext chunkContext) throws Exception { // 从StepExecutionContext获取当前分区的customerId列表 List<String> customerIds = (List<String>) chunkContext.getStepContext() .getStepExecution() .getExecutionContext() .get("partitionedCustomerIds"); // 逐个调用API for (String customerId : customerIds) { apiClient.callCustomerApi(customerId); } return RepeatStatus.FINISHED; } }
- 配置分区Step
把上面的组件整合起来,配置成支持分区的Step,还可以指定异步执行器来并行处理各个分区:
@Bean public Step partitionedCustomerStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, CustomerApiCallTasklet customerApiCallTasklet) { return new StepBuilder("partitionedCustomerStep", jobRepository) .partitioner("customerApiCallStep", new CustomerListPartitioner()) .step(customerApiCallStep(jobRepository, transactionManager, customerApiCallTasklet)) .gridSize(5) // 分区数量,根据你的需求调整 .taskExecutor(new SimpleAsyncTaskExecutor()) // 可选:异步执行分区任务 .build(); } @Bean public Step customerApiCallStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, CustomerApiCallTasklet customerApiCallTasklet) { return new StepBuilder("customerApiCallStep", jobRepository) .tasklet(customerApiCallTasklet, transactionManager) .build(); }
这样每个分区的Tasklet都会拿到对应块大小的customerId列表,完成API调用,和chunk模式的块处理效果一致。
内容的提问来源于stack exchange,提问作者Jazz
相关产品推荐
相关产品推荐

