Spring Batch:未重启JVM时KafkaItemReader重复消费问题排查
你的核心问题是:每次调度启动Spring Batch任务时,KafkaItemReader没有从Kafka的_consumer_offsets主题读取已提交的偏移量,而是重复消费旧数据。这是因为当前配置中,Spring Batch自身的状态管理覆盖了Kafka的偏移量读取逻辑,同时部分Kafka消费者配置也不符合Spring Batch的使用规范。
调整配置的关键点:
关闭Spring Batch的读取器状态保存
当前你设置了saveState(true),这会让KafkaItemReader将偏移量存储到Spring Batch的JobRepository中,每次Job运行时优先读取该存储的状态,而非Kafka的消费者组偏移量。需要将其改为saveState(false),让读取器完全依赖Kafka的偏移量管理:KafkaItemReader<String, byte[]> kafkaItemReader = new KafkaItemReaderBuilder<String, byte[]>() .partitions(partitionsList) .consumerProperties(props) .name("consumer name") .saveState(false) // 关闭状态保存 .topic(topicName) .build();禁用Kafka消费者自动提交偏移量
你的配置中启用了自动提交(ENABLE_AUTO_COMMIT_CONFIG=true),但Spring Batch的KafkaItemReader应该由框架手动提交偏移量,确保偏移量仅在Chunk处理成功后提交。修改该配置为false:props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);移除显式的空偏移量设置
你调用的kafkaItemReader.setPartitionOffsets(new HashMap<>())会强制读取器使用空的初始偏移量配置,这可能干扰Kafka消费者自动获取已提交偏移量的逻辑。直接移除该行代码,读取器会自动根据消费者组ID从_consumer_offsets中获取最新偏移量。
额外验证点:
- 确保所有调度的任务使用同一个group.id:你的配置中已经设置了
GROUP_ID_CONFIG,需保证该值在每次调度时一致,这样Kafka才能识别为同一个消费者组,读取到正确的已提交偏移量。 - 确认Kafka的
auto.offset.reset设置为latest:你的配置中已经正确设置,这会在消费者组无历史偏移量时从最新位置开始消费,避免从头读取。
完成以上调整后,每次调度启动任务时,KafkaItemReader会直接从Kafka消费者组的已提交偏移量开始消费,不再重复处理已完成的记录。
内容的提问来源于stack exchange,提问作者Zied Yazidi

