Spark Streaming将Kafka Avro消息转DataFrame时遇空指针异常求助
解决Spark Streaming读取Kafka Avro消息转DataFrame的空指针问题
咱们先来拆解下你遇到的两个核心问题:
- SparkSession不能在Executor端代码里用:你在DStream的
map算子中调用spark.read.json,但map里的逻辑是跑在Executor节点上的。SparkSession是在Driver端创建的,它没法被序列化传到Executor,所以调用时直接抛出空指针。 - 数据格式不匹配:你的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
相关产品推荐
相关产品推荐

