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

AWS EKS新建Strimzi Kafka集群认证失败求助

AWS EKS中Strimzi Kafka集群Python客户端连接失败排查

问题概述

在AWS EKS中搭建了基于KRaft的Strimzi Kafka集群,尝试通过Python脚本使用SCRAM认证的Kafka用户访问主题时,出现NoBrokersAvailable错误。集群Pod、服务及用户资源均处于正常状态,尝试使用引导服务器LB FQDN和单个Broker LB FQDN均无法解决问题。

集群配置清单

Kafka节点池与集群配置

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaNodePool
metadata:
  name: "kraft-controller"
  labels:
    strimzi.io/cluster: "my-cluster"
spec:
  replicas: 3
  roles:
    - controller
  storage:
    type: jbod
    volumes:
      - id: 0
        type: persistent-claim
        size: "10Gi"
        kraftMetadata: shared
        deleteClaim: false
        class: "test-platform-team-sc"
---
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaNodePool
metadata:
  name: "kafka-broker"
  labels:
    strimzi.io/cluster: "my-cluster"
spec:
  replicas: 3
  roles:
    - broker
  storage:
    type: jbod
    volumes:
      - id: 0
        type: persistent-claim
        size: "10Gi"
        kraftMetadata: shared
        deleteClaim: false
        class: "test-platform-team-sc"
---
apiVersion: kafka.strimzi.io/v1beta2
kind: Kafka
metadata:
  name: "my-cluster"
  annotations:
    strimzi.io/node-pools: enabled
    strimzi.io/kraft: enabled
spec:
  kafka:
    version: "3.8.0"
    metadataVersion: "3.8-IV0"
    authorization:
        type: simple
    listeners:
      - name: plain
        port: 9092
        type: loadbalancer
        tls: false
      - name: tls
        port: 9093
        type: internal
        tls: true
    config:
      offsets.topic.replication.factor: 3
      transaction.state.log.replication.factor: 3
      transaction.state.log.min.isr: 2
      default.replication.factor: 3
      min.insync.replicas: 2
  entityOperator:
    topicOperator:
      watchedNamespace: "strimzi-kafka"
      reconciliationIntervalMs: 60000
      resources:
        requests:
          cpu: "1"
          memory: "500Mi"
        limits:
          cpu: "1"
          memory: "500Mi"
    userOperator:
      watchedNamespace: "strimzi-kafka"
      reconciliationIntervalMs: 60000
      resources:
        requests:
          cpu: "1"
          memory: "500Mi"
        limits:
          cpu: "1"
          memory: "500Mi"

Kafka主题配置

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaTopic
metadata:
  name: "flink-input"
  labels:
    strimzi.io/cluster: "my-cluster"
spec:
  partitions: 10
  replicas: 2
---
apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaTopic
metadata:
  name: "flink-output"
  labels:
    strimzi.io/cluster: "my-cluster"
spec:
  partitions: 10
  replicas: 2

Kafka用户配置

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaUser
metadata:
  name: "flink"
  labels:
    strimzi.io/cluster: "my-cluster"
spec:
  authentication:
    type: scram-sha-512
  authorization:
    type: simple
    acls:
      - resource:
          type: topic
          name: "flink-input"
          patternType: literal
        operations:
          - Create
          - Describe
          - Write
          - Read
        host: "*"
      - resource:
          type: group
          name: "my-group"
          patternType: literal
        operations:
          - Read
        host: "*"
      - resource:
          type: topic
          name: "flink-output"
          patternType: literal
        operations:
          - Create
          - Describe
          - Write
          - Read
        host: "*"

Python测试脚本

import logging
from kafka import KafkaConsumer
from kafka.errors import KafkaError

# Enable logging for kafka-python
logging.basicConfig(level=logging.DEBUG)

# Configuration
BROKER = "AWSELBFQDN:9092"
TOPIC = "flink-input"
USERNAME = "flink"
PASSWORD = "DecodedSecretPassword"

try:
    consumer = KafkaConsumer(
        TOPIC,
        bootstrap_servers=[BROKER],
        security_protocol="SASL_PLAINTEXT",
        sasl_mechanism="PLAIN",
        sasl_plain_username=USERNAME,
        sasl_plain_password=PASSWORD,
        group_id="my-group",
        auto_offset_reset="earliest",
        consumer_timeout_ms=10000
    )

    print("Successfully connected to Kafka broker.")
    print(f"Listening to topic: {TOPIC}")
    for message in consumer:
        print(f"Received message: {message.value.decode('utf-8')}")
        break
    consumer.close()

except KafkaError as e:
    print(f"Failed to connect to Kafka broker: {e}")

错误信息

Failed to connect to Kafka broker: NoBrokersAvailable
pods are healthy and services and user are available

当前集群资源状态

kubectl get pods -n strimzi-kafka
NAME                                          READY   STATUS    RESTARTS   AGE
my-cluster-entity-operator-6f4cc9b654-frzsw   2/2     Running   0          53m
my-cluster-kafka-broker-0                     1/1     Running   0          16m
my-cluster-kafka-broker-1                     1/1     Running   0          15m
my-cluster-kafka-broker-2                     1/1     Running   0          16m
my-cluster-kraft-controller-3                 1/1     Running   0          60m
my-cluster-kraft-controller-4                 1/1     Running   0          61m
my-cluster-kraft-controller-5                 1/1     Running   0          60m
strimzi-cluster-operator-6f9fbb4c75-twx8p     1/1     Running   0          75m

kubectl get services -n strimzi-kafka
NAME                               TYPE           CLUSTER-IP       EXTERNAL-IP                                                                   PORT(S)                               AGE
my-cluster-kafka-bootstrap         ClusterIP      10.100.234.142   <none>                                                                        9091/TCP,9093/TCP                     15d
my-cluster-kafka-broker-plain-0    LoadBalancer   10.100.193.37    awslbfqdn9092:31351/TCP                        17m
my-cluster-kafka-broker-plain-1    LoadBalancer   10.100.234.209   awslbfqdn9092:31518/TCP                        17m
my-cluster-kafka-broker-plain-2    LoadBalancer   10.100.233.6     awslbfqdn9092:31231/TCP                        17m
my-cluster-kafka-brokers           ClusterIP      None             <none>                                                                        9090/TCP,9091/TCP,8443/TCP,9093/TCP   15d
my-cluster-kafka-plain-bootstrap   LoadBalancer   10.100.13.61     awslbfqdn9092:31631/TCP                        17m

kubectl get kafkausers -n strimzi-kafka
NAME    CLUSTER      AUTHENTICATION   AUTHORIZATION   READY
flink   my-cluster   scram-sha-512    simple          True

解决方案建议

1. 修正SASL认证机制

Kafka用户配置的是scram-sha-512认证,但Python脚本中使用的是PLAIN机制,两者不匹配。需将脚本中的sasl_mechanism改为SCRAM-SHA-512。

2. 修正Broker端口配置

从服务列表可见,my-cluster-kafka-plain-bootstrap的外部端口是31631,而非脚本中写的9092。AWS EKS中LoadBalancer服务会将容器端口映射到NodePort,外部访问需使用该NodePort端口。

3. 验证用户密码正确性

确保脚本中的密码是从KafkaUser对应的Secret中正确解码的值,执行以下命令获取密码:

kubectl get secret flink -n strimzi-kafka -o jsonpath='{.data.password}' | base64 -d

4. 检查AWS安全组配置

确保EKS节点的安全组允许客户端IP访问以下端口:

  • Bootstrap LB的NodePort:31631
  • 各Broker LB的NodePort:31351、31518、31231

5. 修改后的Python脚本示例

import logging
from kafka import KafkaConsumer
from kafka.errors import KafkaError

logging.basicConfig(level=logging.DEBUG)

# 修正后的配置
BROKER = "AWSELBFQDN:31631"  # 替换为实际的bootstrap LB外部端口
TOPIC = "flink-input"
USERNAME = "flink"
PASSWORD = "从Secret获取的正确密码"

try:
    consumer = KafkaConsumer(
        TOPIC,
        bootstrap_servers=[BROKER],
        security_protocol="SASL_PLAINTEXT",
        sasl_mechanism="SCRAM-SHA-512",  # 修正为SCRAM-SHA-512
        sasl_plain_username=USERNAME,
        sasl_plain_password=PASSWORD,
        group_id="my-group",
        auto_offset_reset="earliest",
        consumer_timeout_ms=10000
    )

    print("Successfully connected to Kafka broker.")
    print(f"Listening to topic: {TOPIC}")
    for message in consumer:
        print(f"Received message: {message.value.decode('utf-8')}")
        break
    consumer.close()

except KafkaError as e:
    print(f"Failed to connect to Kafka broker: {e}")

6. 验证Kafka advertised地址

若问题仍存在,查看Kafka Broker日志,确认advertised.listeners是否正确设置为LB的外部地址:

kubectl logs my-cluster-kafka-broker-0 -n strimzi-kafka | grep advertised.listeners

内容的提问来源于stack exchange,提问作者Abdul Azizi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 08:54:51