【问题标题】:Kafka connection refused with Kubernetes nodeportKafka 连接被 Kubernetes 节点端口拒绝
【发布时间】:2021-04-18 06:53:43
【问题描述】:

我正在尝试使用节点端口在我的 Kubernetes 设置中公开 KAFKA 以供外部使用。

我的 Helmcharts kafka-service.yaml 如下:

apiVersion: v1
kind: Service
metadata:
  name: kafka
  namespace: test
  labels:
    app: kafka-test
    unit: kafka
spec:
  type: NodePort
  selector:
    app: test-app
    unit: kafka
    parentdeployment: test-kafka
  ports:
    - name: kafka
      port: 9092
      targetPort: 9092
      nodePort: 30092
      protocol: TCP

kafka-deployment.yaml

apiVersion: apps/v1
kind: Deployment
metadata:
  name: kafka
  namespace: {{ .Values.test.namespace }}
  labels:
    app: test-app
    unit: kafka
spec:
  replicas: 1
  template:
    metadata:
      labels:
        app: test-app
        unit: kafka
        parentdeployment: test-kafka
    spec:
      hostname: kafka
      subdomain: kafka
      securityContext:
        fsGroup: {{ .Values.test.groupID }}
      containers:
        - name: kafka
          image: test_kafka:{{ .Values.test.kafkaImageTag }}
          imagePullPolicy: IfNotPresent
          ports:
            - containerPort: 9092
          env:
            - name: IS_KAFKA_CLUSTER
              value: 'false'
            - name: KAFKA_ZOOKEEPER_CONNECT
              value: zookeeper:2281
            - name: KAFKA_LISTENERS
              value: SSL://:9092
            - name: KAFKA_KEYSTORE_PATH
              value: /opt/kafka/conf/kafka.keystore.jks
            - name: KAFKA_TRUSTSTORE_PATH
              value: /opt/kafka/conf/kafka.truststore.jks
            - name: KAFKA_KEYSTORE_PASSWORD
              valueFrom:
                secretKeyRef:
                  name: kafka-secret
                  key: jkskey
            - name: KAFKA_TRUSTSTORE_PASSWORD
              valueFrom:
                secretKeyRef:
                  name: kafka-secret
                  key: jkskey
            - name: KAFKA_LOG_DIRS
              value: /opt/kafka/data
            - name: KAFKA_ADV_LISTENERS
              value: SSL://kafka:9092
            - name: KAFKA_CLIENT_AUTH
              value: none
          volumeMounts:
            - mountPath: "/opt/kafka/conf"
              name: kafka-conf-pv
            - mountPath: "/opt/kafka/data"
              name: kafka-data-pv
      volumes:
        - name: kafka-conf-pv
          persistentVolumeClaim:
            claimName: kafka-conf-pvc
        - name: kafka-data-pv
          persistentVolumeClaim:
            claimName: kafka-data-pvc
  selector:
    matchLabels:
      app: test-app
      unit: kafka
      parentdeployment: test-kafka

动物园管理员服务 yaml

apiVersion: v1
kind: Service
metadata:
  name: zookeeper
  namespace: {{ .Values.test.namespace }}
  labels:
    app: test-ra
    unit: zookeeper
spec:
  type: ClusterIP
  selector:
    app: test-ra
    unit: zookeeper
    parentdeployment: test-zookeeper
  ports:
    - name: zookeeper
      port: 2281

zookeeper 部署 yaml 文件

apiVersion: apps/v1
kind: Deployment
metadata:
  name: zookeeper
  namespace: {{ .Values.test.namespace }}
  labels:
    app: test-app
    unit: zookeeper
spec:
  replicas: 1
  template:
    metadata:
      labels:
        app: test-app
        unit: zookeeper
        parentdeployment: test-zookeeper
    spec:
      hostname: zookeeper
      subdomain: zookeeper
      securityContext:
        fsGroup: {{ .Values.test.groupID }}
      containers:
        - name: zookeeper
          image: test_zookeeper:{{ .Values.test.zookeeperImageTag }}
          imagePullPolicy: IfNotPresent
          ports:
            - containerPort: 2281
          env:
            - name: IS_ZOOKEEPER_CLUSTER
              value: 'false'
            - name: ZOOKEEPER_SSL_CLIENT_PORT
              value: '2281'
            - name: ZOOKEEPER_DATA_DIR
              value: /opt/zookeeper/data
            - name: ZOOKEEPER_DATA_LOG_DIR
              value: /opt/zookeeper/data/log
            - name: ZOOKEEPER_KEYSTORE_PATH
              value: /opt/zookeeper/conf/zookeeper.keystore.jks
            - name: ZOOKEEPER_KEYSTORE_PASSWORD
              valueFrom:
                secretKeyRef:
                  name: zookeeper-secret
                  key: jkskey
            - name: ZOOKEEPER_TRUSTSTORE_PATH
              value: /opt/zookeeper/conf/zookeeper.truststore.jks
            - name: ZOOKEEPER_TRUSTSTORE_PASSWORD
              valueFrom:
                secretKeyRef:
                  name: zookeeper-secret
                  key: jkskey
          volumeMounts:
            - mountPath: "/opt/zookeeper/data"
              name: zookeeper-data-pv
            - mountPath: "/opt/zookeeper/conf"
              name: zookeeper-conf-pv
      volumes:
        - name: zookeeper-data-pv
          persistentVolumeClaim:
            claimName: zookeeper-data-pvc
        - name: zookeeper-conf-pv
          persistentVolumeClaim:
            claimName: zookeeper-conf-pvc
  selector:
    matchLabels:
      app: test-ra
      unit: zookeeper
      parentdeployment: test-zookeeper

kubectl describe for kafka 也显示暴露的节点端口

Type:                     NodePort
IP:                       10.233.1.106
Port:                     kafka  9092/TCP
TargetPort:               9092/TCP
NodePort:                 kafka  30092/TCP
Endpoints:                10.233.66.15:9092
Session Affinity:         None
External Traffic Policy:  Cluster

我有一个发布者二进制文件,它将向 Kafka 发送一些消息。因为我有一个 3 节点集群部署,所以我使用我的主节点 IP 和 Kafka 节点端口 (30092) 来连接 Kafka。

但我的二进制文件出现dial tcp <primary_node_ip>:9092: connect: connection refused 错误。我无法理解为什么即使在 nodePort 到 targetPort 转换成功后它也会被拒绝。随着进一步调试,我在 kafka 日志中看到以下调试日志:

[2021-01-13 08:17:51,692] DEBUG Accepted connection from /10.233.125.0:1564 on /10.233.66.15:9092 and assigned it to processor 0, sendBufferSize [actual|requested]: [102400|102400] recvBufferSize [actual|requested]: [102400|102400] (kafka.network.Acceptor)
[2021-01-13 08:17:51,692] DEBUG Processor 0 listening to new connection from /10.233.125.0:1564 (kafka.network.Processor)
[2021-01-13 08:17:51,702] DEBUG [SslTransportLayer channelId=10.233.66.15:9092-10.233.125.0:1564-245 key=sun.nio.ch.SelectionKeyImpl@43dc2246] SSL peer is not authenticated, returning ANONYMOUS instead (org.apache.kafka.common.network.SslTransportLayer)
[2021-01-13 08:17:51,702] DEBUG [SslTransportLayer channelId=10.233.66.15:9092-10.233.125.0:1564-245 key=sun.nio.ch.SelectionKeyImpl@43dc2246] SSL handshake completed successfully with peerHost '10.233.125.0' peerPort 1564 peerPrincipal 'User:ANONYMOUS' cipherSuite 'TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256' (org.apache.kafka.common.network.SslTransportLayer)
[2021-01-13 08:17:51,702] DEBUG [SocketServer brokerId=1001] Successfully authenticated with /10.233.125.0 (org.apache.kafka.common.network.Selector)
[2021-01-13 08:17:51,707] DEBUG [SocketServer brokerId=1001] Connection with /10.233.125.0 disconnected (org.apache.kafka.common.network.Selector)
java.io.EOFException
        at org.apache.kafka.common.network.SslTransportLayer.read(SslTransportLayer.java:614)
        at org.apache.kafka.common.network.NetworkReceive.readFrom(NetworkReceive.java:95)
        at org.apache.kafka.common.network.KafkaChannel.receive(KafkaChannel.java:448)
        at org.apache.kafka.common.network.KafkaChannel.read(KafkaChannel.java:398)
        at org.apache.kafka.common.network.Selector.attemptRead(Selector.java:678)
        at org.apache.kafka.common.network.Selector.pollSelectionKeys(Selector.java:580)
        at org.apache.kafka.common.network.Selector.poll(Selector.java:485)
        at kafka.network.Processor.poll(SocketServer.scala:861)
        at kafka.network.Processor.run(SocketServer.scala:760)
        at java.lang.Thread.run(Thread.java:748)

使用相同的配置,我能够公开其他服务。我在这里错过了什么?

更新:当我为 EXTERNAL 添加 KAFKA_LISTENERS 和 KAFKA_ADV_LISTENERS 并将 targetPort 更改为 30092 时,外部连接期间的错误消息消失了,但开始出现内部连接的连接错误。

解决方案: 我公开了另一个用于外部通信的服务,如答案中提到的,并将 30092 作为端口和节点端口公开。所以不需要targetPort。我还必须在部署文件中添加额外的 KAFKA_LISTENERS 和 KAFKA_ADV_LISTENERS 以进行外部通信

【问题讨论】:

  • 你能分享一下你的 Kafka Deployment 的 YAML 吗?
  • @EmruzHossain 我已经用 Kafka 部署 yaml 文件更新了这个问题。感谢您的回复。
  • 看起来像一个 ssl 问题。您是否使用某种自签名证书?
  • 我建议考虑使用 Strimzi 在 k8s 上管理 Kafka
  • @ChristophRaab 不。我使用的是 CA 签名证书。

标签: kubernetes apache-kafka kubernetes-helm


【解决方案1】:

我们在一个 Kafka 设置中遇到了类似的问题;我们最终创建了两个 k8s 服务,一个使用 ClusterIP 进行内部通信,第二个具有相同标签的服务使用 NodePort 进行外部通信。

内部访问

apiVersion: v1
kind: Service
metadata:
  name: kafka-internal
  namespace: test
  labels:
    app: kafka-test
    unit: kafka
spec:
  type: NodePort
  selector:
    app: test-app
    unit: kafka
    parentdeployment: test-kafka
  ports:
    - name: kafka
      port: 9092
      protocol: TCP
  type: ClusterIP

外部访问

apiVersion: v1
kind: Service
metadata:
  name: kafka-external
  namespace: test
  labels:
    app: kafka-test
    unit: kafka
spec:
  type: NodePort
  selector:
    app: test-app
    unit: kafka
    parentdeployment: test-kafka
  ports:
    - name: kafka
      port: 9092
      targetPort: 9092
      protocol: TCP
  type: NodePort

【讨论】:

  • 您好,感谢您的回复。我按照您说的尝试了,但收到错误dial tcp: lookup kafka on 169.254.25.10:53: server misbehaving。知道为什么吗?
  • 这是来自您的内部还是外部连接访问?另外,您是否特别想分配 NodePort 30092?如果没有,您可以从外部清单中删除它,并且 k8s 应该分配另一个可用的节点端口。
  • 我需要修复节点端口。错误server misbehaving是因为服务名没有改成kafka-internal。现在我收到dial tcp: lookup kafka-external: Temporary failure in name resolution。
  • 您是否从与 k8s 节点之一相同的主机从二进制文件连接到 kafka-external:30092?还是在单独的工作站上?如果稍后,那么您将不得不使用<NODE_EXTERNAL_IP>:30092
  • 您好,感谢您的回答。我像您提到的那样公开了另一个用于外部通信的服务,并将 30092 公开为它的端口和节点端口。所以不需要targetPort。我还必须在部署文件中添加额外的 KAFKA_LISTENERS 和 KAFKA_ADV_LISTENERS 以进行外部通信
猜你喜欢
  • 2019-05-16
  • 1970-01-01
  • 2019-06-23
  • 2021-09-08
  • 2018-11-26
  • 1970-01-01
  • 2020-09-11
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多