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

Spark Scala应用读取Kafka流AVRO数据反序列化异常求助

问题排查及解决方案

  1. Topic订阅不匹配
    你用kafka-avro-console-consumer消费的Topic是test-py-topic,但Spark代码内配置的订阅Topic为test-2,请先修改option("subscribe", "test-2")为正确的Topic名称。

  2. 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.*")
    
  1. 生产者端数据传入错误
    你的Python AvroProducer生产数据时,传入的value=person.to_json()不符合要求:AvroProducer接收的value应为符合Schema结构的字典对象,传入JSON字符串会导致序列化后的结构和你定义的Avro Schema不匹配,修改为value=person.to_dict()即可。

  2. Schema文件路径拼写错误
    你代码中读取Schema的路径为/src/resources/shemas/schema.avsc,其中shemas为拼写错误,正确应为schemas,请修正路径确保读取到正确的Schema文件。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 05:57:00