使用AvroSerializer推送JSON数据至Kafka Topic失败求助
报错根本原因
你没有正确初始化AvroSerializer实例,错误地将AvroSerializer类本身直接传给了生产者的value.serializer配置项。序列化时SerializingProducer会向序列化器传入(待序列化值, SerializationContext对象)两个参数,此时相当于直接调用AvroSerializer(值, Context对象),而AvroSerializer构造函数第一个参数期望接收字符串类型的Schema,拿到Context对象后调用字符串方法strip()就触发了当前报错。
需修改的4个核心问题
1. 补全缺失的导入
你代码中使用了os.environ但没有导入os模块,开头补充:
import os
2. 修正序列化器初始化逻辑
你用的是AVRO类型Schema,不需要初始化JSONSerializer,直接实例化AvroSerializer,将实例对象传给生产者配置,不要传类本身。同时不需要遍历所有Schema Registry subject,直接按目标subject名获取即可:
if sys.argv[1] == "json" : schema_registry_url = {"url": "http://json_schema-registry:8083"} sr = SchemaRegistryClient(schema_registry_url) # 直接按目标subject名取Schema,不需要遍历所有subject subject = f"{TOPIC}-value" schema = sr.get_latest_version(subject) my_schema = schema.schema.schema_str # 实例化AvroSerializer,是实例不是类 avro_serializer = AvroSerializer(schema_str=my_schema, schema_registry_client=sr) # 配置项传序列化器实例 json_producer_conf = { 'bootstrap.servers':'json_kafka:29094' , 'key.serializer': StringSerializer('utf_8'), 'value.serializer': avro_serializer } producer = SerializingProducer(json_producer_conf)
3. 修正传入生产者的value类型
AvroSerializer接收Python字典类型做序列化,不需要先执行json.dumps转成字符串,修改produce逻辑:
response_json=metrics(1136837) # 直接传字典,不要转JSON字符串 producer.produce(topic=TOPIC, value=response_json) producer.poll(1)
4. 统一Topic命名
你创建Topic时硬编码了"my_topic_in",建议和TOPIC变量统一,避免环境变量配置不一致导致报错:
topic_list = [NewTopic(TOPIC, 1, 1)]
内容的提问来源于stack exchange,提问作者Phil
相关产品推荐
相关产品推荐

