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计数变量没有提前初始化,后续运行也会触发变量未定义错误。
修复步骤
- 修正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重启服务,配置即可生效。 - 如果不想重启服务,可以手动进入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 - 补充代码细节:在循环前初始化计数变量
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
相关产品推荐
相关产品推荐

