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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 22:21:39