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

在Kubernetes中使用KRaft部署Kafka后无法消费消息

Kafka Kubernetes部署消费失败排查问题

我在运行Kubernetes的Docker Desktop环境中,使用以下配置搭建Kafka集群:

apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
  name: kafka-network
spec:
  ingress:
    - from:
        - podSelector:
            matchLabels:
              network/kafka-network: "true"
  podSelector:
    matchLabels:
      network/kafka-network: "true"
---
apiVersion: v1
kind: Service
metadata:
  labels:
    service: kafka-controller
  name: kafka-controller
spec:
  clusterIP: None
  selector:
    app: kafka-controller
    run: kafka-controller
  ports:
    - name: internal
      port: 9093
      targetPort: 9093
---
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: kafka-controller
spec:
  serviceName: kafka-controller
  replicas: 1
  selector:
    matchLabels:
      app: kafka-controller
      run: kafka-controller
  template:
    metadata:
      labels:
        network/kafka-network: "true"
        app: kafka-controller
        run: kafka-controller
    spec:
      hostname: kafka-controller
      securityContext:
        fsGroup: 1000
      enableServiceLinks: false
      containers:
        - name: kafka-controller
          image: apache/kafka:3.9.0
          imagePullPolicy: IfNotPresent
          ports:
            - containerPort: 9093
          env:
            - name: KAFKA_OPTS
              value: "-Djavax.net.debug=all"
            - name: KAFKA_LOG4J_ROOT_LOGLEVEL
              value: debug
            - name: KAFKA_NODE_ID
              value: "0"
            - name: KAFKA_PROCESS_ROLES
              value: "controller"
            - name: KAFKA_LISTENERS
              value: "CONTROLLER://:9093"
            - name: KAFKA_INTER_BROKER_LISTENER_NAME
              value: "PLAINTEXT"
            - name: KAFKA_CONTROLLER_LISTENER_NAMES
              value: "CONTROLLER"
            - name: KAFKA_CONTROLLER_QUORUM_VOTERS
              value: "0@kafka-controller-0.kafka-controller.default.svc.cluster.local:9093"
            - name: KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS
              value: "0"
---
apiVersion: v1
kind: Service
metadata:
  labels:
    service: kafka-broker
  name: kafka-broker
spec:
  clusterIP: None
  selector:
    app: kafka-broker
    run: kafka-broker
  ports:
    - name: internal
      port: 29092
      targetPort: 29092
---
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: kafka-broker
spec:
  serviceName: kafka-broker
  replicas: 1
  selector:
    matchLabels:
      app: kafka-broker
      run: kafka-broker
  template:
    metadata:
      labels:
        network/kafka-network: "true"
        app: kafka-broker
        run: kafka-broker
    spec:
      securityContext:
        fsGroup: 1000
      enableServiceLinks: false
      hostname: kafka-broker
      containers:
        - name: kafka-broker
          image: apache/kafka:3.9.0
          imagePullPolicy: IfNotPresent
          ports:
            - containerPort: 29092
          env:
            - name: KAFKA_OPTS
              value: "-Djavax.net.debug=all"
            - name: KAFKA_NODE_ID
              value: "1"
            - name: KAFKA_PROCESS_ROLES
              value: "broker"
            - name: KAFKA_ADVERTISED_LISTENERS
              value: "PLAINTEXT://kafka-broker-0.kafka-broker.default.svc.cluster.local:29092"
            - name: KAFKA_AUTO_CREATE_TOPICS_ENABLE
              value: "true"
            - name: KAFKA_INTER_BROKER_LISTENER_NAME
              value: "PLAINTEXT"
            - name: KAFKA_CONTROLLER_LISTENER_NAMES
              value: "CONTROLLER"
            - name: KAFKA_LISTENERS
              value: "PLAINTEXT://:29092"
            - name: KAFKA_LISTENER_SECURITY_PROTOCOL_MAP
              value: "PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT"
            - name: KAFKA_CONTROLLER_QUORUM_VOTERS
              value: "0@kafka-controller-0.kafka-controller.default.svc.cluster.local:9093"
            - name: KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS
              value: "0"

执行kubectl apply -f kafka.yml初始化服务后,日志无报错。进入Pod后可正常完成以下操作:

  • 创建topic:
    /opt/kafka/bin/kafka-topics.sh --create --bootstrap-server=kafka-broker:29092 --topic event_topic --partitions 2
    
  • 查看topic信息:
    /opt/kafka/bin/kafka-topics.sh --bootstrap-server=kafka-broker:29092 --describe --topic event_topic
    
  • 生产消息:
    /opt/kafka/bin/kafka-console-producer.sh --bootstrap-server kafka-broker:29092 --topic event_topic
    
  • 通过日志验证消息已写入:
    /opt/kafka/bin/kafka-run-class.sh kafka.tools.DumpLogSegments --deep-iteration --print-data-log --files /tmp/kafka-logs/event_topic-0/00000000000000000000.log
    

但无法通过以下命令消费消息:

/opt/kafka/bin/kafka-console-consumer.sh --topic event_topic --from-beginning --group=c7abf8f3-6ec1-4798-aff5-feeee0886838 --bootstrap-server kafka-broker:29092

已尝试的排查动作:

  • 移除指定消费组后重新执行消费命令
  • 更换bootstrap-server地址为kafka-broker-0.kafka-broker.default.svc.cluster.local:29092或kafka-broker-0:29092

以上操作均无报错,但仍无法消费到消息。

日志观察:

  • Broker日志持续输出:INFO Sent auto-creation request for Set(__consumer_offsets) to the active controller. (kafka.server.DefaultAutoTopicCreationManager)
  • Controller无额外异常日志

补充说明:使用docker-compose部署相同配置的Kafka服务可正常生产消费,怀疑是Kubernetes环境下的配置问题,寻求排查解决思路。


内容的提问来源于stack exchange,提问作者Ernesto Limon

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 10:24:58