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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 03:06:04