无需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格式的消息:
- 修改build.sbt添加依赖:
libraryDependencies += "io.confluent" % "kafka-schema-registry-client" % "7.5.1" libraryDependencies += "io.confluent" % "spark-avro" % "7.5.1" % Provided
- 修改读取代码:
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
相关产品推荐
相关产品推荐

