Spark Scala应用读取Kafka流AVRO数据反序列化异常求助
问题排查及解决方案
Topic订阅不匹配
你用kafka-avro-console-consumer消费的Topic是test-py-topic,但Spark代码内配置的订阅Topic为test-2,请先修改option("subscribe", "test-2")为正确的Topic名称。Confluent Avro序列化前缀适配问题
Confluent Kafka的AvroProducer序列化后的数据会在头部附加5字节前缀(1字节魔术位+4字节Schema ID),Spark原生from_avro函数默认处理无任何前缀的原生Avro数据,直接解析会出现字段为空/默认值的问题,两种修复方案:
- 方案1:手动截断前缀后解析:
val personDF = inputDF.select( from_avro(expr("substring(value, 6)"), jsonFormatSchema).as("person") ).select("person.*") - 方案2:配置Schema Registry地址让Spark自动适配:
val personDF = inputDF.select( from_avro( col("value"), jsonFormatSchema, Map("confluent.schema.registry.url" -> "http://localhost:8081") // 替换为你的Schema Registry实际地址 ).as("person") ).select("person.*")
生产者端数据传入错误
你的Python AvroProducer生产数据时,传入的value=person.to_json()不符合要求:AvroProducer接收的value应为符合Schema结构的字典对象,传入JSON字符串会导致序列化后的结构和你定义的Avro Schema不匹配,修改为value=person.to_dict()即可。Schema文件路径拼写错误
你代码中读取Schema的路径为/src/resources/shemas/schema.avsc,其中shemas为拼写错误,正确应为schemas,请修正路径确保读取到正确的Schema文件。
内容的提问来源于stack exchange,提问作者Illia
相关产品推荐
相关产品推荐

