如何构建SpringBatch架构绕过Oracle IN限制,分批处理采购数据
解决方案:将账户ID读取与采购数据处理合并为单步骤(或使用分区)
你的问题核心在于两个串行步骤会等待第一个步骤完全执行完毕才启动第二个,导致所有账户ID累积后才处理采购数据。要实现每读取1000个账户ID就立即处理对应采购数据,可采用以下两种方案:
方案一:合并为单个Chunk-Oriented步骤
将账户ID读取、采购数据批量查询、Elasticsearch写入整合到一个步骤中,利用Spring Batch的Chunk机制,每处理1000个账户ID(即Chunk Size=1000)就执行一次采购数据的查询与写入,完美适配Oracle的IN查询限制。
1. 账户ID读取器(ItemReader)
从数据库读取账户ID,无需额外累积:
@Bean public ItemReader<Long> accountIdReader(DataSource dataSource) { return new JdbcCursorItemReaderBuilder<Long>() .dataSource(dataSource) .sql("SELECT account_id FROM accounts") .rowMapper((rs, rowNum) -> rs.getLong("account_id")) .build(); }
2. 采购数据写入器(ItemWriter)
接收每一批1000个账户ID,批量查询采购数据后写入ES:
@Bean public ItemWriter<List<Long>> purchaseDataWriter(JdbcTemplate jdbcTemplate, RestHighLevelClient esClient, ObjectMapper objectMapper) { return accountIds -> { // 批量查询采购数据(1000个ID符合Oracle IN查询限制) String sql = "SELECT * FROM purchases WHERE account_id IN (:accountIds)"; MapSqlParameterSource params = new MapSqlParameterSource("accountIds", accountIds); List<Purchase> purchases = jdbcTemplate.query(sql, params, new PurchaseRowMapper()); // 批量写入Elasticsearch BulkRequest bulkRequest = new BulkRequest(); for (Purchase purchase : purchases) { IndexRequest indexRequest = new IndexRequest("purchases") .id(purchase.getId().toString()) .source(objectMapper.writeValueAsString(purchase), XContentType.JSON); bulkRequest.add(indexRequest); } esClient.bulk(bulkRequest, RequestOptions.DEFAULT); }; }
3. 步骤与Job配置
配置Chunk Size为1000,确保每批处理刚好1000个账户ID:
@Bean public Step importPurchaseStep(StepBuilderFactory stepBuilderFactory, ItemReader<Long> accountIdReader, ItemWriter<List<Long>> purchaseDataWriter) { return stepBuilderFactory.get("importPurchaseStep") .<Long, List<Long>>chunk(1000) // 每批处理1000个账户ID .reader(accountIdReader) .writer(purchaseDataWriter) .build(); } @Bean public Job importPurchaseJob(JobBuilderFactory jobBuilderFactory, Step importPurchaseStep) { return jobBuilderFactory.get("importPurchaseJob") .incrementer(new RunIdIncrementer()) .start(importPurchaseStep) .build(); }
方案二:使用分区(Partitioning)并行处理
如果希望并行处理多个批次(提升效率),可以用Spring Batch的分区功能,将10000个账户ID拆分为10个分区(每个1000个),每个分区独立执行采购数据的查询与写入。
1. 分区器(Partitioner)
将账户ID划分为多个分区:
@Bean public Partitioner accountIdPartitioner(JdbcTemplate jdbcTemplate) { return gridSize -> { int totalAccounts = jdbcTemplate.queryForObject("SELECT COUNT(*) FROM accounts", Integer.class); int partitionSize = totalAccounts / gridSize; Map<String, ExecutionContext> partitions = new HashMap<>(); for (int i = 0; i < gridSize; i++) { ExecutionContext context = new ExecutionContext(); int start = i * partitionSize; // 最后一个分区包含剩余所有数据 int end = (i == gridSize - 1) ? totalAccounts : (i + 1) * partitionSize; context.putInt("start", start); context.putInt("end", end); partitions.put("partition-" + i, context); } return partitions; }; }
2. 分区读取器(Step Scope)
根据分区上下文读取对应范围的账户ID:
@Bean @StepScope public ItemReader<Long> accountIdPartitionedReader(DataSource dataSource, @Value("#{stepExecutionContext['start']}") int start, @Value("#{stepExecutionContext['end']}") int end) { return new JdbcCursorItemReaderBuilder<Long>() .dataSource(dataSource) .sql("SELECT account_id FROM accounts ORDER BY account_id OFFSET :start ROWS FETCH NEXT :limit ROWS ONLY") .parameterValues(Map.of("start", start, "limit", end - start)) .rowMapper((rs, rowNum) -> rs.getLong("account_id")) .build(); }
3. 主从步骤与Job配置
// 从步骤:处理单个分区的账户ID @Bean public Step slaveStep(StepBuilderFactory stepBuilderFactory, ItemReader<Long> accountIdPartitionedReader, ItemWriter<List<Long>> purchaseDataWriter) { return stepBuilderFactory.get("slaveStep") .<Long, List<Long>>chunk(1000) .reader(accountIdPartitionedReader) .writer(purchaseDataWriter) .build(); } // 主步骤:协调多个分区并行执行 @Bean public Step masterStep(StepBuilderFactory stepBuilderFactory, Partitioner accountIdPartitioner, Step slaveStep, TaskExecutor taskExecutor) { return stepBuilderFactory.get("masterStep") .partitioner("slaveStep", accountIdPartitioner) .step(slaveStep) .taskExecutor(taskExecutor) // 并行执行的任务执行器 .gridSize(10) // 划分10个分区 .build(); } @Bean public Job importPurchaseJob(JobBuilderFactory jobBuilderFactory, Step masterStep) { return jobBuilderFactory.get("importPurchaseJob") .incrementer(new RunIdIncrementer()) .start(masterStep) .build(); }
关键说明
- 原方案中两个串行步骤的问题:第一个步骤会将所有账户ID处理完毕(例如写入临时存储或内存)后才进入第二个步骤,导致累积全部10000条ID。
- 方案一的优势:简单直接,利用Chunk机制天然实现分批处理,无需额外协调。
- 方案二的优势:适合数据量极大的场景,可通过并行处理提升整体效率。
内容的提问来源于stack exchange,提问作者Jérémie
相关产品推荐
相关产品推荐

