Spark Structured Streaming读取Kafka Avro数据时触发空指针异常
问题:Spark Structured Streaming读取Kafka Avro数据时抛出NullPointerException
代码实现
AvroConsumer 主类
package com.test.spark import com.test.spark.ConfigKafka.getAvroSchema import org.apache.avro.generic.{GenericDatumReader, GenericRecord} import org.apache.avro.io.DecoderFactory import org.apache.spark.SparkContext import org.apache.spark.sql.functions.col import org.apache.spark.sql.{DataFrame, Dataset, Encoder, SparkSession} object AvroConsumer extends App { val sparkSession: SparkSession = SparkSession .builder() .getOrCreate() val sparkContext: SparkContext = sparkSession .sparkContext sparkContext.setLogLevel("WARN") private val topic: String = "topic-test" private val queryName: String = "TEST_QUERY" private val autoOffsetReset: String = "earliest" val avroReader: GenericDatumReader[GenericRecord] = new GenericDatumReader[GenericRecord](getAvroSchema(topic)) val avroDecoderFactory: DecoderFactory = DecoderFactory.get() implicit val encoder: Encoder[GenericRecord] = org.apache.spark.sql.Encoders.kryo[GenericRecord] import sparkSession.implicits._ val kafkaDataFrame: DataFrame = sparkSession .readStream .format("kafka") .option("subscribe", topic) .option("group.id", queryName) .option("startingOffsets", autoOffsetReset) .options(ConfigKafka.getSparkConsumerProperties()) .load() val data: Dataset[String] = kafkaDataFrame .select(col("value").as[Array[Byte]]) .map(d => { val rec: GenericRecord = avroReader.read(null, avroDecoderFactory.binaryDecoder(d, null)) val idContract: String = rec.get("idContract").asInstanceOf[org.apache.avro.util.Utf8].toString println(s"idContract = $idContract") idContract }) data .writeStream .format("console") .option("truncate", false) .start() .awaitTermination() }
ConfigKafka 配置类
object ConfigKafka { def getConfMap(): Map[String, String] = { Map(AbrisConfig.SCHEMA_REGISTRY_URL -> schemaRegistryUrl, SchemaRegistryClientConfig.USER_INFO_CONFIG -> s"${userAccount}:${Encoder.decode(userPassword)}", SchemaRegistryClientConfig.BASIC_AUTH_CREDENTIALS_SOURCE -> "USER_INFO") } def getAvroSchema(inputTopic: String): Schema = { val subjectName: String = s"$inputTopic-value" import collection.JavaConverters._ val confMap: Map[String, String] = getConfMap() val schemaRegistryClient: CachedSchemaRegistryClient = new CachedSchemaRegistryClient(schemaRegistryUrl, 128, confMap.asJava) val stringSchema: String = schemaRegistryClient.getLatestSchemaMetadata(subjectName).getSchema new Schema.Parser().parse(stringSchema) } }
目标Avro Schema
{ "type": "record", "name": "Contract", "namespace": "com.test.data.contract", "fields": [ { "default": null, "name": "entity", "type": [ "null", "string" ] }, { "default": null, "name": "idContract", "type": [ "null", "string" ] }, { "default": null, "name": "codeContract", "type": [ "null", "string" ] } ] }
错误信息
23/05/19 11:16:46 ERROR TaskSetManager: Task 0 in stage 0.0 failed 4 times; aborting job 23/05/19 11:16:46 ERROR WriteToDataSourceV2Exec: Data source write support org.apache.spark.sql.execution.streaming.sources.MicroBatchWrite@5f3442ad is aborting. 23/05/19 11:16:46 ERROR WriteToDataSourceV2Exec: Data source write support org.apache.spark.sql.execution.streaming.sources.MicroBatchWrite@5f3442ad aborted. 23/05/19 11:16:46 ERROR MicroBatchExecution: Query [id = c3424864-6f80-4c29-a316-50949782af69, runId = b4964f51-400b-44a4-8cbc-dc6a16a3922e] terminated with error org.apache.spark.SparkException: Writing job aborted ... Caused by: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 0.0 failed 4 times, most recent failure: Lost task 0.3 in stage 0.0 (TID 3) (host.fr executor 1): java.lang.NullPointerException at com.test.spark.AvroConsumer$.$anonfun$data$1(AvroConsumer.scala:41)
错误指向代码行:
val rec: GenericRecord = avroReader.read(null, avroDecoderFactory.binaryDecoder(d, null))
解决方案
1. 定位NullPointerException的核心原因
你的代码直接解析Kafka消息的原始二进制数据,但大部分Kafka中的Avro消息采用Confluent格式:二进制开头包含5字节头部(1字节魔数 + 4字节Schema ID),直接解析整个字节数组会导致avroReader.read()返回null,进而触发NPE。此外还需考虑:
- Kafka消息的
value字段本身为null - 从Schema Registry获取的Schema与消息实际使用的Schema不匹配
2. 修复手动反序列化逻辑
修改map中的代码,先跳过Confluent Avro的头部,同时增加空值判断:
val data: Dataset[String] = kafkaDataFrame .select(col("value").as[Array[Byte]]) .map(d => { // 处理空消息或格式不合法的消息 if (d == null || d.length < 5) { "" } else { // 跳过5字节的Confluent Avro头部 val payload = d.slice(5, d.length) val rec: GenericRecord = avroReader.read(null, avroDecoderFactory.binaryDecoder(payload, null)) // 用Option处理字段可能为null的情况,避免直接toString触发NPE val idContract = Option(rec.get("idContract")) .map(_.asInstanceOf[org.apache.avro.util.Utf8].toString) .getOrElse("") println(s"idContract = $idContract") idContract } })
3. 使用Abris库简化反序列化(推荐)
你的代码已经引用了Abris相关类,这个库专门用于Spark处理Confluent Avro数据,无需手动处理头部和Schema匹配问题,代码更简洁健壮:
package com.test.spark import za.co.absa.abris.avro.functions.from_avro import za.co.absa.abris.config.AbrisConfig import org.apache.spark.sql.functions.col import org.apache.spark.sql.SparkSession object AvroConsumer extends App { val sparkSession: SparkSession = SparkSession .builder() .getOrCreate() sparkSession.sparkContext.setLogLevel("WARN") private val topic: String = "topic-test" private val queryName: String = "TEST_QUERY" private val autoOffsetReset: String = "earliest" import sparkSession.implicits._ val kafkaDataFrame = sparkSession .readStream .format("kafka") .option("subscribe", topic) .option("group.id", queryName) .option("startingOffsets", autoOffsetReset) .options(ConfigKafka.getSparkConsumerProperties()) .load() // 配置Abris从Schema Registry获取对应Topic的最新Schema val abrisConfig = AbrisConfig .fromConfluentAvro .downloadReaderSchemaByLatestVersion .andTopicNameStrategy(topic) .usingSchemaRegistry(ConfigKafka.getConfMap()) // 直接解析value列,得到结构化DataFrame val parsedDF = kafkaDataFrame .select(from_avro(col("value"), abrisConfig).as("contract")) .select("contract.idContract") parsedDF .writeStream .format("console") .option("truncate", false) .start() .awaitTermination() }
4. 额外注意事项
- 如果消息是用旧版本Schema序列化的,不要直接使用
getLatestSchemaMetadata,可以指定Schema ID或版本号获取对应Schema - 生产环境中避免使用
println,改用Spark的日志系统输出 - 处理流式数据时,建议增加异常捕获逻辑,避免单个坏消息导致整个流任务终止
内容的提问来源于stack exchange,提问作者Mamaf
相关产品推荐
相关产品推荐

