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

Spark Streaming读取Kafka遇空指针异常:存在空Offset数据

问题分析与解决方案

问题根源

从错误栈可以看出,NullPointerException发生在KafkaRecordToUnsafeRowConverter.toUnsafeRow方法中——这说明Spark在将Kafka原始记录转换为内部UnsafeRow格式时就崩溃了,早于你在foreachBatch中执行的过滤逻辑。所以你在foreachBatch里写的filter(col("value").isNotNull)根本没机会运行,自然无法解决问题。

结合你提供的Kafka控制台输出,问题出现在Kafka中存在value为null的消息(比如Debezium同步MongoDB删除事件时可能产生这类消息),而Spark的Kafka数据源在处理null类型的value字段时触发了底层的空指针异常。

解决方案

方案1:在流转换链中提前处理/过滤null值

在load()之后立即对value字段进行处理,避免后续转换时触发NPE。可以用coalesce将null的value替换为合法的空二进制数据,再进行过滤:

import org.apache.spark.sql.functions.{col, coalesce, length}

val kafkaDataFrameRaw = spark
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", kafkaUrl)
    .option("subscribe", topics)
    .option("maxOffsetsPerTrigger", maxOffsetsPerTrigger)
    .option("startingOffsets", startOffsetFinal)
    .option("failOnDataLoss", false)
    .load()
    // 将null的value替换为空二进制,避免转换崩溃
    .selectExpr(
      "key",
      "coalesce(value, cast('' as binary)) as value",
      "topic",
      "partition",
      "offset",
      "timestamp",
      "timestampType"
    )
    // 过滤掉空的value(包括替换后的空二进制)
    .filter(length(col("value")) > 0)

kafkaDataFrameRaw.foreachBatch { (result: DataFrame, batchId: Long) =>
    result.show(10,false) 
}.start()

方案2:使用Spark Avro处理(适配你的Avro场景)

既然你用了kafka-avro-console-consumer,说明消息是Avro格式。可以结合spark-avro库解析消息,同时提前过滤null值:

import org.apache.spark.sql.avro._
import org.apache.spark.sql.functions.col

// 从Schema Registry获取最新的Avro schema
val schema = spark.read
  .format("avro")
  .option("avroSchemaRegistryUrl", "http://schemaregurl01:8081")
  .option("avroSchemaSubject", s"$topics-value")
  .load()
  .schema

val kafkaDataFrameRaw = spark
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", kafkaUrl)
    .option("subscribe", topics)
    .option("maxOffsetsPerTrigger", maxOffsetsPerTrigger)
    .option("startingOffsets", startOffsetFinal)
    .option("failOnDataLoss", false)
    .load()
    // 先过滤掉value为null的记录
    .filter(col("value").isNotNull)
    // 解析Avro格式消息
    .select(from_avro(col("value"), schema).alias("data"))

kafkaDataFrameRaw.foreachBatch { (result: DataFrame, batchId: Long) =>
    result.show(10,false) 
}.start()

方案3:使用低级API手动处理消息

如果上述方案仍无法解决,可以切换到Spark Streaming低级API(DirectStream),直接控制每条消息的读取与过滤:

import org.apache.spark.streaming.{StreamingContext, Seconds}
import org.apache.spark.streaming.kafka010._
import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe

// 初始化StreamingContext
val ssc = new StreamingContext(spark.sparkContext, Seconds(5))

val kafkaParams = Map[String, Object](
  "bootstrap.servers" -> kafkaUrl,
  "key.deserializer" -> classOf[org.apache.kafka.common.serialization.StringDeserializer],
  "value.deserializer" -> classOf[org.apache.kafka.common.serialization.StringDeserializer],
  "group.id" -> "your-spark-group-id",
  "auto.offset.reset" -> "latest",
  "enable.auto.commit" -> (false: java.lang.Boolean)
)

val topicsSet = topics.split(",").toSet
val kafkaStream = KafkaUtils.createDirectStream[String, String](
  ssc,
  PreferConsistent,
  Subscribe[String, String](topicsSet, kafkaParams)
)

// 直接过滤掉value为null的记录
val filteredStream = kafkaStream.filter(_.value() != null)

// 后续处理逻辑
filteredStream.foreachRDD { rdd =>
  val df = spark.read.json(rdd.map(_.value()))
  df.show(10, false)
}

ssc.start()
ssc.awaitTermination()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 10:43:09