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}-valuesubject下,需要确认你的账号有对应注册权限- 如果你需要提前校验本地schema的合法性,可以单独用
avro.schema.Parse做校验,但不要将解析后的对象传给AvroSerializer
内容的提问来源于stack exchange,提问作者Eypros
相关产品推荐
相关产品推荐

