You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Cloud Stream Kafka Binder从指定offset获取消息及监听器无响应问题

监听器无法触发的原因

  1. 偏移量指向分区末尾无可用消息
    从运行日志可以明确看到Resetting offset for partition MY_TOPIC-0 to offset 1076,这里的1076是该分区当前的最大偏移量(LEO,代表下一条待写入消息的偏移量),该分区已有的1076条消息偏移量范围为0~1075,你从1076位置开始消费,自然没有存量消息可以拉取,没有新消息写入的前提下监听器不会被触发。
  2. 配置层级错误导致参数不生效
    你的application.yml中kafka、cloud配置直接顶格编写,没有放在Spring Boot标准的spring父节点下,导致你配置的消费组test-kafka-service等参数完全不生效,日志中显示的消费组ID为latest就是默认配置的结果。
  3. 日志打印语法错误
    监听方法中的日志编写为log.info("*** MESSAGE: ***", msg),缺少占位符{},就算有消息进入方法也无法打印出消息内容,会误以为监听器没有触发,正确写法为log.info("*** MESSAGE: {} ***", msg)。
  4. 手动提交偏移量未实现
    你配置了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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 01:36:03