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

本地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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 22:46:03