本地K8s集群KRaft模式Kafka本地Python脚本生产Avro消息方案
问题:本地Python脚本连接KRaft模式Kafka(Docker Desktop K8s集群)的最简配置方案
背景
我通过Tiltfile在本地Docker Desktop K8s集群部署了以下服务:
- Kafka Broker(镜像
confluentinc/confluent-local:7.5.0,默认运行于KRaft模式) - Kafka Schema Registry(镜像
confluentinc/cp-schema-registry:7.5.2) - 包含消费者逻辑的MyApp服务
我需要用本地Python脚本向Kafka Broker发送消息,测试MyApp的消费者逻辑。此前在集群内通过shell进入Kafka Pod执行kafka-console-producer的方式存在问题:
- 主题采用Avro Schema,
kafka-console-producer发送的消息会报无效魔术字节错误 - 重复操作效率低下,必须改用本地脚本
最初尝试通过kubectl port-forward暴露Kafka Broker,让Python脚本连接localhost:<port>,但参考的旧文章仅适用于ZooKeeper模式,按其配置修改后Kafka Pod出现CONTROLLER相关错误,复杂配置均未生效。现寻求最简配置方案,附上最后可用的Tiltfile配置及尝试后失效的配置:
最后可用的Tiltfile配置
deployment_create( name='kafka-deployment', image='artifactory.alteryx.com/docker/confluentinc/confluent-local:7.5.0', ports=['9092'], env=[ {'name': 'KAFKA_SASL_MECHANISM', 'value': 'PLAIN' }, {'name': 'KAFKA_ADVERTISED_LISTENERS', 'value': 'PLAINTEXT://kafka-deployment:9092,PLAINTEXT_HOST://kafka-deployment:9092'}, {'name': 'KAFKA_SCHEMA_REGISTRY_URL', 'value': "http://schema-registry-deployment:8081"}, {'name': 'KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE', 'value': 'true'}, {'name': 'KAFKA_LOG_CLEANUP_POLICY', 'value': 'compact'} ] ) deployment_create( name='schema-registry-deployment', image='artifactory.alteryx.com/docker/confluentinc/cp-schema-registry:7.5.2', ports=['8081'], env=[ {'name': 'SCHEMA_REGISTRY_HOST_NAME', 'value': 'schema-registry-deployment' }, {'name': 'SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS', 'value': 'PLAINTEXT://kafka-deployment:9092' } ], deps=['kafka-deployment'] )
尝试后失效的Tiltfile配置
deployment_create( name='kafka-deployment', image='confluentinc/confluent-local:7.5.0', ports=['9092'], env=[ {'name': 'KAFKA_SASL_MECHANISM', 'value': 'PLAIN' }, {'name': 'KAFKA_LISTENERS', 'value': 'INTERNAL_LISTENER://kafka-deployment:29092, EXTERNAL_LISTENER://localhost:9092, CONTROLLER://0.0.0.0:29093'}, {'name': 'KAFKA_ADVERTISED_LISTENERS', 'value': 'INTERNAL_LISTENER://kafka-deployment:29092,EXTERNAL_LISTENER://localhost:9092'}, {'name': 'KAFKA_SCHEMA_REGISTRY_URL', 'value': "http://schema-registry-deployment:8081"}, {'name': 'KAFKA_LISTENER_SECURITY_PROTOCOL_MAP', 'value': 'EXTERNAL_LISTENER:PLAINTEXT,INTERNAL_LISTENER:PLAINTEXT,CONTROLLER:PLAINTEXT'}, {'name': 'KAFKA_INTER_BROKER_LISTENER_NAME', 'value': 'INTERNAL_LISTENER'}, {'name': 'KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE', 'value': 'true'}, {'name': 'KAFKA_LOG_CLEANUP_POLICY', 'value': 'compact'} ] ) deployment_create( name='schema-registry-deployment', image='confluentinc/cp-schema-registry:7.5.2', ports=['8081'], env=[ {'name': 'SCHEMA_REGISTRY_HOST_NAME', 'value': 'schema-registry-deployment'}, {'name': 'SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS', 'value': 'INTERNAL_LISTENER://kafka-deployment:29092'} ], deps=['kafka-deployment'] )
内容的提问来源于stack exchange,提问作者omgirok
相关产品推荐
相关产品推荐

