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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 21:07:46