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
相关产品推荐
相关产品推荐

