【问题标题】:Unable to connect kafka in kubernetes cluster from python-producer无法从 python-producer 连接 kubernetes 集群中的 kafka
【发布时间】:2022-10-24 10:21:48
【问题描述】:

我在一个 kubernetes 集群中有一个 python fastapi 应用程序和 kafka,但在不同的命名空间中。 Python 应用程序生成从客户端到 kafka 主题的消息。首先,在 python 脚本中,我尝试通过以下方式连接到 kafka 代理:

producer = KafkaProducer(
    bootstrap_servers=os.environ["KAFKA_SERVER"],
    sasl_plain_username=os.environ["KAFKA_BROKER_USERNAME"],
    sasl_plain_password=os.environ["KAFKA_BROKER_PASSWORD"],
    security_protocol="PLAINTEXT",
    sasl_mechanism="PLAINTEXT",
    value_serializer=lambda v: v.encode('utf-8')
)

KAFKA_SERVER 的值是集群中 kafka 服务的名称。在这种情况下是:gb-kafka.kafka.svc.gb.local:9092 当应用程序启动时,它会崩溃并引发错误:

Traceback (most recent call last):
  File "/code/main.py", line 45, in <module>
    producer = KafkaProducer(
  File "/usr/local/lib/python3.10/site-packages/kafka/producer/kafka.py", line 381, in __init__
    client = KafkaClient(metrics=self._metrics, metric_group_prefix='producer',
  File "/usr/local/lib/python3.10/site-packages/kafka/client_async.py", line 244, in __init__
    self.config['api_version'] = self.check_version(timeout=check_timeout)
  File "/usr/local/lib/python3.10/site-packages/kafka/client_async.py", line 927, in check_version
    raise Errors.NoBrokersAvailable()
kafka.errors.NoBrokersAvailable: NoBrokersAvailable 

我在 Helm 的帮助下,准确地使用 bitnami/kafka 配置了 kafka。我没有对 values.yaml 文件做一些大的改变,只是改变了监听器的配置:

listeners: "PLAINTEXT://:9092"
advertisedListeners: "PLAINTEXT://gb-kafka.kafka.svc.gb.local:9092"
listenerSecurityProtocolMap: "PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT"
allowPlaintextListener: true
interBrokerListenerName: PLAINTEXT

kubectl get all 命令得到以下结果(ip 地址已更改):

NAME                       READY   STATUS    RESTARTS   AGE
pod/gb-kafka-0             1/1     Running   0          45m
pod/gb-kafka-zookeeper-0   1/1     Running   0          4d2h

NAME                                  TYPE        CLUSTER-IP      EXTERNAL-IP   PORT(S)                      AGE
service/gb-kafka                      ClusterIP   10.10.10.543    <none>        9092/TCP                     4d2h
service/gb-kafka-headless             ClusterIP   None            <none>        9092/TCP,9093/TCP            4d2h
service/gb-kafka-zookeeper            ClusterIP   10.222.01.220   <none>        2181/TCP,2888/TCP,3888/TCP   4d2h
service/gb-kafka-zookeeper-headless   ClusterIP   None            <none>        2181/TCP,2888/TCP,3888/TCP   4d2h

NAME                                  READY   AGE
statefulset.apps/gb-kafka             1/1     4d2h
statefulset.apps/gb-kafka-zookeeper   1/1     4d2h

最有趣的是,如果我只是运行简单的 pod,使用 python 连接到具有相同配置和简单生产者的 kafka,而不使用 fastapi,它可以很好地连接到代理。

【问题讨论】:

  • 如果您使用纯文本侦听器,则不需要 SASL 属性。此外,不再维护 kafka python 包本身
  • 但在不同的命名空间- 您可能需要一个 NetworkPolicy 来允许不同命名空间中的 pod 相互访问

标签: kubernetes apache-kafka pykafka


【解决方案1】:

如果您使用的是 bitnami/kafka 图表,则从 kubernetes 集群内部连接时无需更改侦听器配置。

您只需通过 kubernetes 的内部 dns 策略访问 kafka。

尝试gb-kafka.kafka.svc.cluster.local:9092或者gb-kafka-headless.kafka.svc.cluster.local:9092

【讨论】:

    猜你喜欢
    • 2019-06-10
    • 2021-03-10
    • 1970-01-01
    • 2020-06-10
    • 2018-12-28
    • 2020-08-15
    • 2019-07-31
    • 2020-09-11
    • 1970-01-01
    相关资源
    最近更新 更多