为什么KafkaItemReader在新作业执行时总是包含上一次运行的最后一条偏移量记录?
问题原因及解决方案
核心原因
你当前配置的KafkaItemReader是默认单例Bean,多个定时触发的作业实例会复用同一个Reader对象:
- 上一次作业执行完成后,Reader内部缓存的偏移量状态没有被重置
- 新作业启动时不会重新从Spring Batch的JobRepository中读取上次持久化的最终偏移量,直接复用了缓存的旧偏移量,就会重复读取上一次最后一条处理的记录
- 应用重启时会重新初始化Reader Bean,会主动从持久化状态加载偏移量,所以首次运行正常
解决方案
给KafkaItemReader添加@StepScope注解即可:
@Bean @StepScope public KafkaItemReader<String, String> kafkaItemReader() { Properties props = new Properties(); // 原有配置保持不变,可额外添加以下配置避免自动提交偏移量冲突 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 其余原有逻辑保持不变 }
@StepScope会让每次Step执行时都创建全新的Reader实例,不会复用旧的状态- 新实例初始化时会自动从JobRepository读取上一次成功执行的最终偏移量,从偏移量+1的位置开始读取,不会出现重复消费问题
额外校验点
你当前的作业参数用时间作为唯一标识、作业配置了RunIdIncrementer、saveState=true的配置都没有问题,不需要额外调整。
内容的提问来源于stack exchange,提问作者vick
相关产品推荐
相关产品推荐

