Spring Cloud Stream Kafka Binder从指定offset获取消息及监听器无响应问题
监听器无法触发的原因
- 偏移量指向分区末尾无可用消息
从运行日志可以明确看到Resetting offset for partition MY_TOPIC-0 to offset 1076,这里的1076是该分区当前的最大偏移量(LEO,代表下一条待写入消息的偏移量),该分区已有的1076条消息偏移量范围为0~1075,你从1076位置开始消费,自然没有存量消息可以拉取,没有新消息写入的前提下监听器不会被触发。 - 配置层级错误导致参数不生效
你的application.yml中kafka、cloud配置直接顶格编写,没有放在Spring Boot标准的spring父节点下,导致你配置的消费组test-kafka-service等参数完全不生效,日志中显示的消费组ID为latest就是默认配置的结果。 - 日志打印语法错误
监听方法中的日志编写为log.info("*** MESSAGE: ***", msg),缺少占位符{},就算有消息进入方法也无法打印出消息内容,会误以为监听器没有触发,正确写法为log.info("*** MESSAGE: {} ***", msg)。 - 手动提交偏移量未实现
你配置了autoCommitOffset: false关闭了自动偏移量提交,但没有在监听方法中编写手动提交偏移量的逻辑,即使有消息消费成功,偏移量也不会更新,后续消费也会出现异常。
实现从指定offset读取消息的方案
方案1:配置固定起始偏移量(推荐固定场景使用)
先修正yml配置的层级,然后添加偏移量重置相关参数即可,示例配置如下:
spring: kafka: consumer: properties: max.poll.interval.ms: 3600000 max-poll-records: 10 cloud: zookeeper: connect-string: test.kafka.com:2181,test.kafka.com:2181,test.kafka.com:2181 stream: kafka: bindings: my-group-id: consumer: autoCommitOffset: false resetOffsets: true startOffset: 0 # 此处填写你需要的起始偏移量,也可填earliest(最早)、latest(最新) binder: brokers: - test.kafka.com:6667 - test.kafka.com:6667 - test.kafka.com:6667 auto-create-topics: false auto-add-partitions: false jaas: controlFlag: REQUIRED loginModule: com.sun.security.auth.module.Krb5LoginModule options: useKeyTab: true storeKey: true serviceName: kafka keyTab: C:\\files\\user.keytab principal: user@test.com debug: true configuration: security: protocol: SASL_PLAINTEXT bindings: my-group-id: binder: kafka destination: MY_TOPIC group: test-kafka-service servlet: multipart: max-file-size: 50MB max-request-size: 50MB
注意:该配置仅在消费组没有已提交的偏移量时生效,若要强制重置可以更换新的消费组ID测试。
方案2:代码动态指定偏移量(推荐灵活调整场景使用)
可以在监听方法中获取Consumer对象,手动调用seek方法指定偏移量,示例代码如下:
@Slf4j @Component @RequiredArgsConstructor @EnableBinding(EventConsumer.class) public class EventListener { @StreamListener(target = "my-group-id") public void processMessage(Object msg, @Header(KafkaHeaders.CONSUMER) Consumer<?, ?> consumer) { // 首次消费时调用,指定MY_TOPIC的0分区从偏移量100开始消费 consumer.seek(new TopicPartition("MY_TOPIC", 0), 100); log.info("*** MESSAGE: {} ***", msg); // 消息处理完成后手动提交偏移量 consumer.commitSync(); } }
内容的提问来源于stack exchange,提问作者Java Student
相关产品推荐
相关产品推荐

