【发布时间】: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 属性。此外,不再维护
kafkapython 包本身 -
但在不同的命名空间- 您可能需要一个 NetworkPolicy 来允许不同命名空间中的 pod 相互访问
标签: kubernetes apache-kafka pykafka