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

无需Abris或移除魔术字节,Spark Scala反序列化Confluent Kafka Avro消息

问题:Spark反序列化Confluent Avro Kafka消息失败(非Abris/魔术字节移除方案)

需求与环境

  • 目标:反序列化Confluent Avro格式的Kafka消息,不使用Abris库或手动移除魔术字节
  • 环境:Spark 3.1.2,Scala 2.12.14

报错信息

Caused by: org.apache.avro.AvroRuntimeException: Malformed data. Length is negative: -5057
at org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:308)

尝试代码

import java.nio.file.{Files, Paths}
import org.apache.kafka.clients.consumer.ConsumerConfig
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.avro.functions._
import org.apache.spark.sql.streaming.Trigger
import io.confluent.kafka.serializers.KafkaAvroDeserializer

import org.apache.avro.Schema
import org.apache.avro.Schema.Parser
import org.apache.http.client.methods.HttpGet
import org.apache.http.impl.client.HttpClients
import org.apache.http.util.EntityUtils

object ReadFromKafka {
  System.setProperty("hadoop.home.dir", "C:\\hadoop")
  def main(args: Array[String]): Unit = {
    val sparkConf = new SparkConf()
    sparkConf.setAppName("ConfluentConsumer").setMaster("local[*]")
    val spark = SparkSession.builder()
      .appName("AvroKafkaStructuredStreaming")
      .config(sparkConf)
      .getOrCreate()

    import spark.implicits._
    val kafkaBootstrapServers = "bootstrapserver"
    val kafkaTopic = "topic"
    val schemaRegistryUrl = "schemaregistryurl:8081"

    // Function to fetch Avro schema from Schema Registry
    def fetchAvroSchema(topic: String, schemaRegistryUrl: String, version: String): Schema = {
      val client = HttpClients.createDefault()
      val url = s"$schemaRegistryUrl/subjects/$topic-value/versions/$version/schema"
      println("url--"+url)
      val request = new HttpGet(url)
      val response = client.execute(request)
      val jsonSchema = EntityUtils.toString(response.getEntity)
      client.close()
      new Parser().parse(jsonSchema)
    }

    val df = spark.read
      .format("kafka")
      .option("kafka.bootstrap.servers", kafkaBootstrapServers)
      .option("subscribe", kafkaTopic)
      .option("startingOffsets", "earliest")
      .option("failOnDataLoss", "false")
      .option("maxOffsetsPerTrigger", 10)
      .option("value.deserializer", classOf[KafkaAvroDeserializer].getName)
      .option("schema.registry.url", schemaRegistryUrl)
      .option("mode", "PERMISSIVE")
      .load()

    import org.apache.kafka.clients.consumer.ConsumerConfig

    // Fetch Avro schema dynamically
    val avroSchema = fetchAvroSchema(kafkaTopic, schemaRegistryUrl, "latest").toString()
    println("avroSchema "+avroSchema)
    
    val deserializedDF = df.select(from_avro($"value", avroSchema).as("data"))
    // val deserializedDF =  df.select(from_avro(expr("substring(value, 6)"), avroSchema)) // this works

    deserializedDF.show(false)

    /* val query = deserializedDF.writeStream
       .outputMode("append")
       .format("console")
     //  .trigger(Trigger.ProcessingTime("10 seconds"))
       .start()*/

    //query.awaitTermination()
  }
}

build.sbt配置

scalaVersion := "2.12.14"
val sparkVersion = "3.1.2"

resolvers += "confluent" at "https://packages.confluent.io/maven/"
resolvers += "confluent" at "https://packages.confluent.io/maven/"
libraryDependencies += "org.apache.spark" %% "spark-sql-kafka-0-10" % "3.1.2"
libraryDependencies += "io.confluent" % "kafka-avro-serializer" % "7.5.1"
libraryDependencies += "org.apache.spark" %% "spark-avro" % "3.1.2"
libraryDependencies ++= Seq(
  "org.apache.spark" %% "spark-core" % sparkVersion %"compile",
  "org.apache.spark" %% "spark-sql" % sparkVersion % "compile"
)

dependencyOverrides ++= {
  Seq(
    "com.fasterxml.jackson.module" %% "jackson-module-scala" % "2.12.3",
    "com.fasterxml.jackson.core" % "jackson-annotations" % "2.12.3",
    "com.fasterxml.jackson.core" % "jackson-core" % "2.12.3",
    "com.fasterxml.jackson.core" % "jackson-databind" % "2.12.3"
  )
}

问题原因

Spark Structured Streaming的Kafka数据源不支持value.deserializer参数,该参数仅适用于原生Kafka Consumer。实际读取时,消息值会以原始二进制(BinaryType)返回,KafkaAvroDeserializer根本没有生效。而from_avro函数只能解析标准Avro二进制格式,无法识别Confluent Avro特有的前缀(1字节魔术位+4字节Schema ID),直接解析就会出现格式错误。

可行解决方案(不依赖Abris/手动删魔术字节)

方案1:自定义UDF使用Confluent KafkaAvroDeserializer

通过自定义UDF直接调用Confluent的反序列化器处理二进制消息,无需手动处理魔术字节:

import org.apache.kafka.common.serialization.Deserializer
import io.confluent.kafka.serializers.KafkaAvroDeserializer
import org.apache.spark.sql.functions.udf

// 初始化Confluent Avro反序列化器(单例模式避免重复初始化)
lazy val avroDeserializer: Deserializer[Any] = {
  val deserializer = new KafkaAvroDeserializer()
  deserializer.configure(
    Map(
      "schema.registry.url" -> schemaRegistryUrl,
      "specific.avro.reader" -> "false" // 使用通用Avro记录,无需生成具体类
    ).asJava,
    false
  )
  deserializer
}

// 定义反序列化UDF
val deserializeConfluentAvro = udf((value: Array[Byte]) => {
  if (value == null) null else avroDeserializer.deserialize(null, value)
})

// 应用UDF解析消息
val deserializedDF = df.select(deserializeConfluentAvro($"value").as("data"))
deserializedDF.show(false)

注意:如果是流处理场景,确保反序列化器是可序列化的(将其放在object中或使用lazy val)。

方案2:使用Confluent官方Spark Avro集成

Confluent提供了专门适配Structured Streaming的Avro数据源,可直接读取Confluent格式的消息:

  1. 修改build.sbt添加依赖:
libraryDependencies += "io.confluent" % "kafka-schema-registry-client" % "7.5.1"
libraryDependencies += "io.confluent" % "spark-avro" % "7.5.1" % Provided
  1. 修改读取代码:
val df = spark.read
  .format("io.confluent.kafka.serializers.AbstractKafkaAvroSerDeConfig")
  .option("kafka.bootstrap.servers", kafkaBootstrapServers)
  .option("subscribe", kafkaTopic)
  .option("startingOffsets", "earliest")
  .option("schema.registry.url", schemaRegistryUrl)
  .load()

df.show(false)

注意:Confluent版本需与Spark版本兼容,Spark 3.1.2可搭配Confluent 7.x系列。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 04:45:57