【问题标题】:NullPointerException in SQLContext.read() SparkSQLContext.read() Spark 中的 NullPointerException
【发布时间】:2016-08-10 11:13:36
【问题描述】:

我正在尝试使用 SQLContext.read() 在 Spark 中读取由 Kafka 生成的 JSON 记录。每次出现 NullPointerException。

    SparkConf conf = new SparkConf()
       .setAppName("kafka-sandbox")
        .setMaster("local[*]");
    JavaSparkContext sc = new JavaSparkContext(conf);
    JavaStreamingContext ssc = new JavaStreamingContext(sc, new Duration(2000));

    Set<String> topics = Collections.singleton(topicString);
    Map<String, String> kafkaParams = new HashMap<>();
    kafkaParams.put("metadata.broker.list", servers);

    JavaPairInputDStream<String, String> directKafkaStream = KafkaUtils.createDirectStream(
            ssc, String.class, String.class, StringDecoder.class, StringDecoder.class, 
            kafkaParams, topics);
    SQLContext sqlContext = new org.apache.spark.sql.SQLContext(sc);

    directKafkaStream
        .map(message -> message._2)
        .foreachRDD(rdd -> {
            rdd.foreach(record -> {
                Dataset<Row> ds = sqlContext.read().json(rdd);
            });
         });
    ssc.start();
    ssc.awaitTermination();

这是一个日志:

java.lang.NullPointerException
    at org.apache.spark.sql.SparkSession.sessionState$lzycompute(SparkSession.scala:112)
    at org.apache.spark.sql.SparkSession.sessionState(SparkSession.scala:110)
    at org.apache.spark.sql.DataFrameReader.<init>(DataFrameReader.scala:535)
    at org.apache.spark.sql.SparkSession.read(SparkSession.scala:595)
    at org.apache.spark.sql.SQLContext.read(SQLContext.scala:504)
    at SparkJSONConsumer$1.lambda$2(SparkJSONConsumer.java:73)
    at SparkJSONConsumer$1$$Lambda$8/1821075039.call(Unknown Source)
    at org.apache.spark.api.java.JavaRDDLike$$anonfun$foreach$1.apply(JavaRDDLike.scala:350)
    at org.apache.spark.api.java.JavaRDDLike$$anonfun$foreach$1.apply(JavaRDDLike.scala:350)
    at scala.collection.Iterator$class.foreach(Iterator.scala:893)
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1336)
    at org.apache.spark.rdd.RDD$$anonfun$foreach$1$$anonfun$apply$27.apply(RDD.scala:875)
    at org.apache.spark.rdd.RDD$$anonfun$foreach$1$$anonfun$apply$27.apply(RDD.scala:875)
    at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1897)
    at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1897)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:70)
    at org.apache.spark.scheduler.Task.run(Task.scala:85)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:274)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    at java.lang.Thread.run(Thread.java:745)

我认为问题是由于 foreachRDD 子句引起的,但无法弄清楚。所以任何建议都会很棒。

另外,我正在使用 sqlContext,因为在我计划以 avro 格式(“com.databricks.spark.avro”)序列化记录之后。如果有办法将包含JSON结构的字符串序列化为avro格式而不定义schema,非常欢迎分享!

提前致谢。

【问题讨论】:

    标签: java json apache-spark avro


    【解决方案1】:

    如 Spark 文档中所述 - 您必须使用 StreamingContext 正在使用的 SparkContext 创建一个 SparkSession。此外,必须这样做,以便它可以在驱动程序故障时重新启动。这是通过创建一个惰性实例化的 SparkSession 单例实例来完成的。

    参考:

    http://spark.apache.org/docs/2.1.0/streaming-programming-guide.html#dataframe-and-sql-operations

    解决方案:

    在读取 json 之前创建如下 SQLContext。

    SQLContext sqlContext = SparkSession.builder.config(rdd.sparkContext.getConf).getOrCreate().sqlContext
    
    
    Dataset<Row> ds = sqlContext.read().json(rdd);
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2017-03-17
      • 1970-01-01
      • 1970-01-01
      • 2017-05-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多