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

Python向Kafka双分区发送数据报Unrecognized partition错误

报错根因

报错本质是代码中指定的目标分区,在生产者拉取到的Kafka主题元数据中不存在,具体触发原因有两个:

  • 主题名不匹配:生产者代码往topic_test2发送消息,但docker-compose配置里的KAFKA_CREATE_TOPICS参数仅配置了自动创建topic_test(2分区1副本),没有自动创建业务实际要用的topic_test2。
  • 自动创建的主题分区数不符合要求:Kafka默认开启自动创建主题配置,第一次往topic_test2发第一条消息(指定partition=0)时,Kafka会按默认参数创建主题,默认分区数为1,即仅存在分区0,不存在分区1。当第二条消息指定发往分区1时,生产者从本地缓存的元数据中查不到对应分区,就会抛出Unrecognized partition断言错误。
    另外代码中循环使用的aux计数变量没有提前初始化,后续运行也会触发变量未定义错误。
修复步骤
  1. 修正Docker Compose配置,把自动创建的主题名改成实际使用的topic_test2,指定2分区1副本:
    environment:
      # 其余原有配置保持不变,仅修改下面这行
      KAFKA_CREATE_TOPICS: "topic_test2:2:1"
    
    修改后执行docker-compose -f docker-compose-expose.yml down再docker-compose -f docker-compose-expose.yml up -d重启服务,配置即可生效。
  2. 如果不想重启服务,可以手动进入Kafka容器创建符合配置要求的主题:
    # 替换为实际的kafka容器ID/名称,可通过docker ps查询
    docker exec -it <kafka容器ID/名称> bash
    # 先删除之前错误自动创建的1分区topic_test2,从未创建过该主题可跳过这步
    kafka-topics.sh --delete --topic topic_test2 --zookeeper zookeeper:2181
    # 创建2分区1副本的目标主题
    kafka-topics.sh --create --topic topic_test2 --partitions 2 --replication-factor 1 --zookeeper zookeeper:2181
    # 执行命令确认主题配置正确,返回结果中能看到2个分区即配置生效
    kafka-topics.sh --describe --topic topic_test2 --zookeeper zookeeper:2181
    
  3. 补充代码细节:在循环前初始化计数变量aux = 0;如果启动生产者后立刻发消息偶发识别不到分区,可以在发消息前主动拉取一次主题元数据,参考代码如下:
    import json
    from kafka import KafkaProducer
    producer = KafkaProducer(bootstrap_servers='localhost:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8'))
    # 主动拉取主题元数据,阻塞直到拿到完整元数据
    producer.partitions_for('topic_test2')
    
    aux = 0
    for tweet in hashtag_tweets:
        data = {}
        data["content"] = tweet._json["full_text"]
        data["retweets"] = tweet._json["retweet_count"]
        data["favorites"] = tweet._json["favorite_count"]
        if (aux % 2 == 0):
            producer.send('topic_test2', value=data, partition=0)
        else:
            producer.send('topic_test2', value=data, partition=1)  
        aux  += 1
    producer.flush()
    producer.close()
    

注意:如果之前已经自动创建了1分区的topic_test2,必须先删除旧主题再创建新主题,不然分区数修改不会生效。

内容的提问来源于stack exchange,提问作者María Jesús Cañoles

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 20:54:25