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

Spark Streaming将Kafka Avro消息转DataFrame时遇空指针异常求助

解决Spark Streaming读取Kafka Avro消息转DataFrame的空指针问题

咱们先来拆解下你遇到的两个核心问题:

  1. SparkSession不能在Executor端代码里用:你在DStream的map算子中调用spark.read.json,但map里的逻辑是跑在Executor节点上的。SparkSession是在Driver端创建的,它没法被序列化传到Executor,所以调用时直接抛出空指针。
  2. 数据格式不匹配:你的Kafka消息是Avro格式,却用read.json解析,这本身就不对,就算没有空指针,解析也会失败。

下面给你两种解决方案,优先推荐用Structured Streaming(Spark 2.x之后的主流API,比DStream更简洁易维护),也会给你DStream的修正方案。

方案一:改用Structured Streaming(推荐)

Structured Streaming原生支持Kafka和Avro,处理结构化数据更顺畅:

package consumerTest
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._

object Consumer {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder
      .master("local")
      .appName("my-spark-app")
      .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
      .getOrCreate()

    import spark.implicits._

    // Kafka配置
    val kafkaParams = Map[String, String](
      "bootstrap.servers" -> "<kafka-server>:9092",
      "group.id" -> "sakwq",
      "auto.offset.reset" -> "earliest",
      "enable.auto.commit" -> "false"
    )

    val topic = "cdcemployee"

    // 从Kafka读取消息
    val kafkaDF = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", kafkaParams("bootstrap.servers"))
      .option("subscribe", topic)
      .option("group.id", kafkaParams("group.id"))
      .option("auto.offset.reset", kafkaParams("auto.offset.reset"))
      .option("enable.auto.commit", kafkaParams("enable.auto.commit"))
      .load()

    // 解析Avro格式的value,直接从Schema Registry拉取最新 schema
    val avroDF = kafkaDF.select(
      from_avro(
        col("value"), 
        "http://<schema-registry>:8181/subjects/cdcemployee-value/versions/latest"
      ).as("employee_data")
    ).select("employee_data.*") // 展开嵌套结构

    // 将数据写入内存临时表,也可以换成其他存储(比如Hive、JDBC)
    val query = avroDF.writeStream
      .format("memory")
      .queryName("employee_table")
      .outputMode("append")
      .start()

    query.awaitTermination()
  }
}

依赖说明(以SBT为例)

需要引入这些依赖才能正常运行:

libraryDependencies ++= Seq(
  "org.apache.spark" %% "spark-sql" % "2.4.8", // 对应你的Spark版本
  "org.apache.spark" %% "spark-sql-kafka-0-10" % "2.4.8",
  "org.apache.spark" %% "spark-avro" % "2.4.8",
  "io.confluent" % "kafka-avro-serializer" % "5.4.0" // 对应你的Confluent版本
)

方案二:修正DStream代码

如果你一定要用DStream,核心是不能在Executor的map里用SparkSession,改用foreachRDD在Driver端处理每个RDD,同时正确解析Avro:

package consumerTest
import io.confluent.kafka.serializers.KafkaAvroDeserializer
import org.apache.spark.sql.{SparkSession}
import org.apache.spark.streaming._
import org.apache.spark.streaming.kafka010._
import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe
import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
import org.apache.avro.generic.GenericRecord

object Consumer {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder
      .master("local")
      .appName("my-spark-app")
      .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
      .getOrCreate()

    import spark.implicits._
    val ssc = new StreamingContext(spark.sparkContext, Seconds(10))

    val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> "<kafka-server>:9092",
      "key.deserializer" -> classOf[KafkaAvroDeserializer],
      "value.deserializer" -> classOf[KafkaAvroDeserializer],
      "group.id" -> "sakwq",
      "auto.offset.reset" -> "earliest",
      "enable.auto.commit" -> "false",
      "schema.registry.url" -> "http://<schema-registry>:8181"
    )

    val topics = Array("cdcemployee")
    // 把value类型指定为GenericRecord,方便操作Avro字段
    val stream = KafkaUtils.createDirectStream[String, GenericRecord](
      ssc,
      PreferConsistent,
      Subscribe[String, GenericRecord](topics, kafkaParams)
    )

    // 用foreachRDD在Driver端处理每个批次的RDD
    stream.foreachRDD { rdd =>
      // 转换RDD为DataFrame,这里需要根据你的Avro字段调整
      val employeeDF = rdd.map(_.value())
        .map { avroRecord =>
          (
            avroRecord.get("id").asInstanceOf[Int],
            avroRecord.get("name").asInstanceOf[String],
            avroRecord.get("department").asInstanceOf[String]
          )
        }.toDF("id", "name", "department")

      // 注册临时表
      employeeDF.createOrReplaceTempView("employee_table")

      // 测试查询临时表
      spark.sql("SELECT * FROM employee_table LIMIT 10").show()
    }

    ssc.start()
    ssc.awaitTermination()
  }
}

关键注意点

  • 把Object换成GenericRecord,可以直接访问Avro的字段
  • foreachRDD的逻辑是在Driver端触发的,所以可以安全使用SparkSession
  • 生产环境要去掉master("local"),换成集群模式的配置

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:03:11