单分区Kafka Topic消息逐个消费配置缺失问题排查
问题根源
你之所以会遇到这个情况,是因为Spring Cloud Stream Kafka Binder的消费者默认配置在悄悄“帮”你拉取多条消息:默认的max.poll.records值是500,也就是说消费者会一次性从Kafka的这个单分区拉取最多500条消息到本地缓存。哪怕你还没对当前消息做偏移量确认,框架也会把本地缓存里的下一条消息推送给你的监听方法——这就造成了“未确认就处理下一条”的假象,本质是这些消息已经提前被拉到本地了。
解决办法
你只需要在消费者配置里加上maxPollRecords: 1,强制消费者每次只拉取一条消息。这样只有当你调用acknowledgment.acknowledge()确认当前消息的偏移量后,消费者才会去拉取下一条消息,完全符合你“逐个同步处理”的需求。
修改后的完整配置如下:
spring: application: name: file-consumer cloud: stream: kafka: binder: type: kafka brokers: localhost defaultBrokerPort: 29092 defaultZkPort: 32181 configuration: max.request.size: 300000 max.message.bytes: 300000 bindings: fileWriteBindingInput: consumer: autoCommitOffset: false # 添加这行,限制每次仅拉取1条消息 maxPollRecords: 1 bindings: fileWriteBindingInput: binder: kafka destination: files.write group: ${spring.application.name} contentType: 'text/plain'
额外说明
可能你之前的理解是“未提交偏移量就不会消费下一条”,这个逻辑其实没错,但要注意:Kafka的消费者是先拉取消息到本地缓存,再逐步处理。如果一次拉取了多条,那本地缓存里的消息会被依次处理,和偏移量是否提交无关。只有当本地缓存的消息都处理完后,下次poll才会根据已提交的偏移量去拉新的消息。所以设置max.poll.records=1就能从根源上保证每次只处理一条,处理完确认后再拉下一条。
内容的提问来源于stack exchange,提问作者OlivierTerrien
相关产品推荐
相关产品推荐

