在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
相关产品推荐
相关产品推荐

