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

confluent-kafka读取本地Avro schema并将消息key序列化为字符串方法

错误原因

报错是因为你传入了avro.schema.Parse返回的RecordSchema对象给AvroSerializer,或者手动处理schema字符串时破坏了原有合法的schema结构。confluent-kafka的AvroSerializer要求schema_str参数传入原生的JSON格式schema字符串,不需要提前用avro库做解析,也不需要手动移除空格、换行符,这类操作可能会损坏schema中包含空格/换行的字段说明、默认值等内容,导致schema不合法。

正确实现代码

你只需要直接读取本地.avsc文件的原始内容传入即可,修改后的代码段如下:

# 读取本地schema文件,直接获取原始字符串即可
with open(args.schema, "r", encoding="utf-8") as f:
    schema_str = f.read()

pro_conf = {"auto.register.schemas": True}
# 直接传入原始schema字符串,不需要用avro.schema.Parse解析
avro_serializer = AvroSerializer(
    schema_registry_client=schema_registry_client,
    schema_str=schema_str,
    conf=pro_conf
)
注意事项
  • 如果不需要提前校验本地schema合法性,完全不需要导入avro库做解析操作
  • auto.register.schemas设为True时,Schema Registry会自动将本地schema注册到对应主题的{topic}-value subject下,需要确认你的账号有对应注册权限
  • 如果你需要提前校验本地schema的合法性,可以单独用avro.schema.Parse做校验,但不要将解析后的对象传给AvroSerializer

内容的提问来源于stack exchange,提问作者Eypros

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 00:15:09