KafkaJS无法连接Kubernetes Pod内的Kafka Broker问题
问题概述
在Kubernetes Pod中部署Kafka后,使用Node.js的KafkaJS连接Broker时,无论是用kafka-srv:9092还是kafka-srv.default.svc.cluster.local:9092,都出现getaddrinfo ENOTFOUND的DNS解析错误,无法建立连接。
错误日志
{"level":"ERROR","timestamp":"2024-06-01T15:18:02.594Z","logger":"kafkajs","message":"[Connection] Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local","broker":"kafka-srv.default.svc.cluster.local:9092","clientId":"ticketing","stack":"Error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local\n at GetAddrInfoReqWrap.onlookupall [as oncomplete] (node:dns:118:26)"} {"level":"ERROR","timestamp":"2024-06-01T15:18:02.595Z","logger":"kafkajs","message":"[BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local","retryCount":3,"retryTime":2242} {"level":"ERROR","timestamp":"2024-06-01T15:18:04.915Z","logger":"kafkajs","message":"[Connection] Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local","broker":"kafka-srv.default.svc.cluster.local:9092","clientId":"ticketing","stack":"Error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local\n at GetAddrInfoReqWrap.onlookupall [as oncomplete] (node:dns:118:26)"} {"level":"ERROR","timestamp":"2024-06-01T15:18:04.916Z","logger":"kafkajs","message":"[BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local","retryCount":4,"retryTime":5076} {"level":"ERROR","timestamp":"2024-06-01T15:18:10.054Z","logger":"kafkajs","message":"[Connection] Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local","broker":"kafka-srv.default.svc.cluster.local:9092","clientId":"ticketing","stack":"Error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local\n at GetAddrInfoReqWrap.onlookupall [as oncomplete] (node:dns:118:26)"} {"level":"ERROR","timestamp":"2024-06-01T15:18:10.055Z","logger":"kafkajs","message":"[BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local","retryCount":5,"retryTime":11724} Error starting the consumer: KafkaJSNonRetriableError Caused by: KafkaJSConnectionError: Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local at Socket.onError (C:\Users\almog\Downloads\programming\programming\microservices\ticketing\kafka\node_modules\kafkajs\src\network\connection.js:210:23) ... 3 lines matching cause stack trace ... at processTicksAndRejections (node:internal/process/task_queues:82:21) { name: 'KafkaJSNumberOfRetriesExceeded', retriable: false, helpUrl: undefined, retryCount: 5, retryTime: 11724, [cause]: KafkaJSConnectionError: Connection error: getaddrinfo ENOTFOUND kafka-srv.default.svc.cluster.local at Socket.onError (C:\Users\almog\Downloads\programming\programming\microservices\ticketing\kafka\node_modules\kafkajs\src\network\connection.js:210:23) at Socket.emit (node:events:518:28) at emitErrorNT (node:internal/streams/destroy:169:8) at emitErrorCloseNT (node:internal/streams/destroy:128:3) at processTicksAndRejections (node:internal/process/task_queues:82:21) { retriable: true, helpUrl: undefined, broker: 'kafka-srv.default.svc.cluster.local:9092', code: 'ENOTFOUND', [cause]: undefined } }
Kubernetes部署配置(kafka-depl.yaml)
apiVersion: apps/v1 kind: Deployment metadata: name: kafka-depl spec: replicas: 1 selector: matchLabels: app: kafka template: metadata: labels: app: kafka spec: containers: - name: kafka image: bitnami/kafka:latest ports: - containerPort: 9092 resources: requests: memory: "1Gi" cpu: "500m" limits: memory: "2Gi" cpu: "1" env: - name: KAFKA_BROKER_ID value: "1" - name: KAFKA_LISTENERS value: PLAINTEXT://:9092 - name: KAFKA_ADVERTISED_LISTENERS value: PLAINTEXT://kafka-srv:9092 - name: KAFKA_ZOOKEEPER_CONNECT value: zookeeper-srv:2181 - name: ALLOW_PLAINTEXT_LISTENER value: "yes" - name: KAFKA_CFG_LISTENERS value: PLAINTEXT://:9092 - name: KAFKA_CFG_ADVERTISED_LISTENERS value: PLAINTEXT://kafka-srv:9092 - name: KAFKA_CFG_ZOOKEEPER_CONNECT value: zookeeper-srv:2181 --- apiVersion: v1 kind: Service metadata: name: kafka-srv spec: selector: app: kafka ports: - name: kafka protocol: TCP port: 9092 targetPort: 9092 --- apiVersion: apps/v1 kind: Deployment metadata: name: zookeeper-depl spec: replicas: 1 selector: matchLabels: app: zookeeper template: metadata: labels: app: zookeeper spec: containers: - name: zookeeper image: bitnami/zookeeper:latest ports: - containerPort: 2181 resources: requests: memory: "512Mi" cpu: "250m" limits: memory: "1Gi" cpu: "500m" env: - name: ALLOW_ANONYMOUS_LOGIN value: "yes" --- apiVersion: v1 kind: Service metadata: name: zookeeper-srv spec: selector: app: zookeeper ports: - name: zookeeper protocol: TCP port: 2181 targetPort: 2181
消费者代码(consumer.ts)
import { Kafka } from "kafkajs"; import { TicketCreatedConsumer } from "./events/ticket-created-consumer"; console.clear(); const kafka = new Kafka({ clientId: "ticketing", brokers: ["kafka-srv:9092"], }); const run = async () => { const ticketCreatedConsumer = new TicketCreatedConsumer(kafka); await ticketCreatedConsumer.consume(); const closeKafka = async () => { await ticketCreatedConsumer.close(); process.exit(); }; process.on("SIGINT", closeKafka); process.on("SIGTERM", closeKafka); console.log("Consumer connected to Kafka"); }; run().catch((err) => { console.error("Error starting the consumer:", err); });
解决方案
1. 确认消费者Pod与Kafka服务同命名空间
Kubernetes服务的短域名(如kafka-srv)仅在同一命名空间内可解析。如果消费者Pod部署在不同命名空间,必须使用完整域名kafka-srv.<命名空间>.svc.cluster.local,或者调整两者到同一命名空间。
2. 验证Kafka服务的端点状态
执行以下命令检查Kafka服务是否正确关联到Pod:
kubectl get endpoints kafka-srv
输出应显示Kafka Pod的IP和端口9092。如果端点为空,说明Deployment的标签与Service的selector不匹配,检查两者的app标签是否一致(当前配置中均为app: kafka,需确认实际部署后标签未被修改)。
3. 检查消费者Pod的DNS解析能力
进入消费者Pod内部,测试对Kafka服务的DNS解析:
kubectl exec -it <consumer-pod-name> -- nslookup kafka-srv # 或测试完整域名 kubectl exec -it <consumer-pod-name> -- nslookup kafka-srv.default.svc.cluster.local
如果解析失败,说明Kubernetes DNS组件(如CoreDNS)存在问题,检查DNS Pod状态:
kubectl get pods -n kube-system | grep coredns
若DNS Pod异常,重启或重建CoreDNS部署。
4. 修正Kafka环境变量冗余配置
Bitnami Kafka镜像优先使用带_CFG_前缀的环境变量,当前配置中重复设置了KAFKA_LISTENERS和KAFKA_CFG_LISTENERS等变量,可能导致配置冲突。删除冗余的非_CFG_变量,保留以下配置即可:
env: - name: KAFKA_BROKER_ID value: "1" - name: KAFKA_CFG_LISTENERS value: PLAINTEXT://:9092 - name: KAFKA_CFG_ADVERTISED_LISTENERS value: PLAINTEXT://kafka-srv:9092 - name: KAFKA_CFG_ZOOKEEPER_CONNECT value: zookeeper-srv:2181 - name: ALLOW_PLAINTEXT_LISTENER value: "yes"
修改后重新部署Kafka:
kubectl apply -f kafka-depl.yaml
5. 确认消费者代码的Broker地址
若消费者与Kafka在同一命名空间,保持brokers: ["kafka-srv:9092"]即可;若在不同命名空间,替换为对应完整域名。
内容的提问来源于stack exchange,提问作者Almog Hindi

