Spring Batch每日定时分块处理:如何周期性发送100条数据?
解决方案
你的核心问题在于ItemReader一次性加载了所有数据到内存,导致后续的chunk分块只是在内存中拆分,而非真正的读100条发100条、循环执行。要实现每读取100条就立刻发送,需要从ItemReader的实现和任务执行逻辑两方面调整:
1. 实现分页式ItemReader
使用Spring Batch提供的JdbcPagingItemReader(适配MySQL等关系型数据库),它支持按页读取数据,每次只加载指定数量的记录(即你的chunk size=100),读完一页再自动读取下一页,直到没有数据为止。
示例配置:
@Bean public ItemReader<Order> itemReader(DataSource dataSource) { JdbcPagingItemReader<Order> reader = new JdbcPagingItemReader<>(); reader.setDataSource(dataSource); reader.setPageSize(100); // 每次读取100条,和chunk size一致 reader.setRowMapper(new BeanPropertyRowMapper<>(Order.class)); // 映射数据库记录到Order对象 // 配置分页查询规则(以MySQL为例) MySqlPagingQueryProvider queryProvider = new MySqlPagingQueryProvider(); queryProvider.setSelectClause("SELECT id, order_no, user_id, product_id, create_time"); queryProvider.setFromClause("FROM orders"); queryProvider.setWhereClause("status = 'UNPROCESSED'"); // 只读取未处理的数据,避免重复发送 // 指定排序键,保证分页读取的一致性(必须有唯一且稳定的排序字段,比如id) queryProvider.setSortKeys(Collections.singletonMap("id", Order.ASCENDING)); reader.setQueryProvider(queryProvider); return reader; }
2. 完善ItemWriter:发送Kafka+标记已处理
发送Kafka消息后,需要批量更新数据库中对应数据的状态,避免每日任务重复处理相同记录:
@Bean public ItemWriter<Order> orderKafkaSender(KafkaTemplate<String, Order> kafkaTemplate, JdbcTemplate jdbcTemplate) { return orders -> { // 批量发送Kafka消息 orders.forEach(order -> kafkaTemplate.send("your-kafka-topic", order.getId().toString(), order)); // 批量更新数据状态为已处理 List<Long> orderIds = orders.stream().map(Order::getId).collect(Collectors.toList()); String updateSql = "UPDATE orders SET status = 'PROCESSED' WHERE id IN (:ids)"; MapSqlParameterSource parameters = new MapSqlParameterSource().addValue("ids", orderIds); jdbcTemplate.update(updateSql, parameters); }; }
3. 保持原Step/Job配置(无需修改核心逻辑)
你的Step和Job配置本身是正确的,只要ItemReader是分页实现,chunk(100)就会自动按“读100条→发100条→读下100条”的逻辑循环:
@Bean public Step sendUsersOrderProductsStep(ItemReader<Order> itemReader, ItemWriter<Order> orderKafkaSender) throws Exception { return this.stepBuilderFactory.get("testStep") .<Order, Order>chunk(100) .reader(itemReader) .writer(orderKafkaSender) .build(); } @Bean Job sendOrdersJob(Step sendUsersOrderProductsStep) throws Exception { return this.jobBuilderFactory.get("testJob") .start(sendUsersOrderProductsStep) .build(); }
4. 配置每日定时任务
使用Spring的@Scheduled实现每日触发,注意要给每次Job执行添加唯一参数(避免Spring Batch判定为重复任务):
@Configuration @EnableScheduling public class BatchScheduler { @Autowired private JobLauncher jobLauncher; @Autowired private Job sendOrdersJob; // 每日凌晨1点执行,可根据需求调整cron表达式 @Scheduled(cron = "0 0 1 * * ?") public void runDailyOrderJob() throws Exception { JobParameters jobParameters = new JobParametersBuilder() .addString("run.timestamp", String.valueOf(System.currentTimeMillis())) .toJobParameters(); jobLauncher.run(sendOrdersJob, jobParameters); } }
关键注意事项
- 排序键必须唯一且稳定:分页查询依赖排序键保证每次读取的分页数据不重复、不遗漏,建议用数据库主键(如id)。
- 避免重复处理:通过status字段标记已处理数据,是保证每日任务正确性的核心,否则每次任务都会重新发送所有历史数据。
- Job参数唯一性:Spring Batch不允许同一Job使用相同参数重复执行,所以每次触发都要添加唯一标识(如时间戳)。
内容的提问来源于stack exchange,提问作者kamal adel
相关产品推荐
相关产品推荐

