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

KafkaJS无法连接Kubernetes Pod内的Kafka Broker问题

Kubernetes中Kafka DNS解析失败问题排查与解决

问题概述

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 04:42:07