Kafka Avro反序列化Python报错:默认值类型不匹配求助
问题描述
我是Kafka和Python的新手,需要创建Kafka消费者。已经实现简单消费者并能获取结果,但Kafka中存储的是Avro格式数据,因此需要进行反序列化。尝试了如下代码:
import os from confluent_kafka import Consumer from confluent_kafka.serialization import SerializationContext, MessageField from confluent_kafka.schema_registry import SchemaRegistryClient from confluent_kafka.schema_registry.avro import AvroDeserializer if __name__ == "__main__": class test(object): def __init__(self,test_id=None,dep=None,descr=None,stor_key=None,pos=None,time_dt=None): self.test_id = test_id self.dep = dep self.descr = descr self.stor_key = stor_key self.pos = pos self.time_dt = time_dt def dict_to_klf(obj, ctx): if obj is None: return None return test(test_id=obj['test_id'], dep=obj['dep'], descr=obj['descr'], stor_key=obj['stor_key'], pos=obj['pos'], time_dt=obj['time_dt']) schema = "descr.avsc" path = os.path.realpath(os.path.dirname(__file__)) with open(f"{path}\\{schema}") as f: schema_str = f.read() sr_conf = {'url': ':8081'} schema_registry_client = SchemaRegistryClient(sr_conf) avro_deserializer = AvroDeserializer(schema_registry_client, schema_str, dict_to_klf) consumer_config = { "bootstrap.servers": "com:9092", "group.id": "descr_events", "auto.offset.reset": "earliest" } consumer = Consumer(consumer_config) consumer.subscribe(['descr']) while True: msg = consumer.poll(1) if msg is None: continue user = avro_deserializer(message.value, SerializationContext(topic, MessageField.VALUE)) print(msg.topic()) print("-------------------------")
运行后出现报错:
fastavro._schema_common.SchemaParseException: Default value <undefined> must match schema type: long
对应的Avro Schema文件descr.avsc内容如下:
{ "type": "record", "name": "klf", "namespace": "test_ns", "fields": [ { "name": "descr", "type": "string", "default": "undefined" }, { "name": "test_id", "type": "long", "default": "undefined" }, { "name": "dep", "type": "string", "default": "undefined" }, { "name": "stor_key", "type": "string", "default": "undefined" }, { "name": "time_dt", "type": "string", "default": "undefined" }, { "name": "pos", "type": "string", "default": "undefined" } ] }
需要修改哪些内容才能正常获取数据?
解决方法
1. 修正Avro Schema的默认值类型错误
报错核心原因是test_id字段类型为long,但默认值设置为字符串"undefined",类型不匹配。Avro要求字段默认值必须和字段类型一致,有两种修正方案:
方案1:设置数字类型的默认值
将test_id的默认值改为符合long类型的数字,比如0:
{ "type": "record", "name": "klf", "namespace": "test_ns", "fields": [ { "name": "descr", "type": "string", "default": "undefined" }, { "name": "test_id", "type": "long", "default": 0 }, { "name": "dep", "type": "string", "default": "undefined" }, { "name": "stor_key", "type": "string", "default": "undefined" }, { "name": "time_dt", "type": "string", "default": "undefined" }, { "name": "pos", "type": "string", "default": "undefined" } ] }
方案2:允许字段为null
修改test_id的类型为联合类型["long", "null"],并将默认值设为null:
{ "type": "record", "name": "klf", "namespace": "test_ns", "fields": [ { "name": "descr", "type": "string", "default": "undefined" }, { "name": "test_id", "type": ["long", "null"], "default": null }, { "name": "dep", "type": "string", "default": "undefined" }, { "name": "stor_key", "type": "string", "default": "undefined" }, { "name": "time_dt", "type": "string", "default": "undefined" }, { "name": "pos", "type": "string", "default": "undefined" } ] }
2. 修正消费者代码中的变量错误
代码存在两处变量名称错误,同时建议增加错误处理:
while True: msg = consumer.poll(1) if msg is None: continue # 增加错误判断 if msg.error(): print(f"Consumer error: {msg.error()}") continue # 修正变量名:message改为msg,调用value()方法;topic改为msg.topic() user = avro_deserializer(msg.value(), SerializationContext(msg.topic(), MessageField.VALUE)) print(msg.topic()) # 可选:打印反序列化后的对象内容 if user: print(f"test_id: {user.test_id}, descr: {user.descr}") print("-------------------------")
3. 补充Schema Registry地址
当前sr_conf中的url为空,需要填写实际的Schema Registry服务地址,比如:
sr_conf = {'url': 'http://your-schema-registry-host:8081'}
内容的提问来源于stack exchange,提问作者Denys
相关产品推荐
相关产品推荐

