【问题标题】:1.5.1| Spark Streaming | NullPointerException with SQL createDataFrame1.5.1|火花流 |带有 SQL createDataFrame 的 NullPointerException
【发布时间】:2016-07-22 00:24:15
【问题描述】:

我使用的是 Spark 1.5.1。

在流context 中,我得到SQLContext 如下

SQLContext sqlContext = SQLContext.getOrCreate(records.context()); DataFrame dataFrame = sqlContext.createDataFrame(record, SchemaRecord.class); dataFrame.registerTempTable("records");

records 是一个 JavaRDD 每个 Record 都有以下结构

public class SchemaRecord implements Serializable {

private static final long serialVersionUID = 1L; 
private String msisdn; 
private String application_type; 
//private long uplink_bytes = 0L;
}

当 msisdn 和 application_type 等字段类型只是 strings 时,一切正常。

当我添加 Long 类型的 uplink_bytes 之类的另一个字段时,我得到了 在 createDataFrame 出现 NullPointer 异常

Exception in thread "main" java.lang.NullPointerException
at org.spark-project.guava.reflect.TypeToken.method(TypeToken.java:465)
at 
org.apache.spark.sql.catalyst.JavaTypeInference$$anonfun$2.apply(JavaTypeInference.scala:103)
at 
org.apache.spark.sql.catalyst.JavaTypeInference$$anonfun$2.apply(JavaTypeInference.scala:102)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:244)
at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:244)
at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:108)
at scala.collection.TraversableLike$class.map(TraversableLike.scala:244)
at scala.collection.mutable.ArrayOps$ofRef.map(ArrayOps.scala:108)
at org.apache.spark.sql.
catalyst.JavaTypeInference$.org$apache$spark$sql$catalyst$JavaTypeInference$$inferDataType(JavaTypeInference.scala:102)
at org.apache.spark.sql.catalyst.JavaTypeInference$.inferDataType(JavaTypeInference.scala:47)
at org.apache.spark.sql.SQLContext.getSchema(SQLContext.scala:1031)
at org.apache.spark.sql.SQLContext.createDataFrame(SQLContext.scala:519)
at org.apache.spark.sql.SQLContext.createDataFrame(SQLContext.scala:548)

请推荐

【问题讨论】:

  • 你是如何创建 DataFrame 的?动态还是手动?能否请您发布完整的代码。您是否还将“SchemaRecord”定义为案例类,以防它是动态的?
  • @Sumit - 请找到已编辑的问题以回复您的问题
  • @Sumit 请原谅我不太明白。 Schema Rechord 是一个普通的 Java 对象。不是动态的。
  • 尝试使用 Long(包装类)而不是“long”。那应该行得通。还要确保您的 RDD 确实包含 Long 类型值。
  • @Sumit,早些时候它实际上是 Long 类型而不是原始类型,然后我尝试改变一些东西以使它们工作,但它们都没有:(

标签: apache-spark apache-spark-sql spark-streaming


【解决方案1】:

您的问题可能是您的模型类不是一个干净的 JavaBean。目前 Spark 没有代码来处理具有 setter 但没有 getter 方法的属性。你可以简单地尝试这样的事情来检查 Spark 如何理解你的类:

PropertyDescriptor[] props = Introspector.getBeanInfo(YourClass.class).getPropertyDescriptors();
for(PropertyDescriptor prop:props) {
    System.out.println(prop.getDisplayName());
    System.out.println("\t"+prop.getReadMethod());
    System.out.println("\t"+prop.getWriteMethod());
}

内省器还将只有 setter 的字段识别为预操作,这会在 Spark 中引发 NullPointerException。

【讨论】:

  • 从 spark documentation 开始,Spark SQL 支持自动将 JavaBeans 的 RDD 转换为 DataFrame。使用反射获得的 BeanInfo 定义了表的模式。目前,Spark SQL 不支持包含嵌套或复杂类型(如列表或数组)的 JavaBean。您可以通过创建一个实现 Serializable 并为其所有字段具有 getter 和 setter 的类来创建 JavaBean。
【解决方案2】:

这是我尝试过的,它奏效了:-

这里是存储 String、Long 和 Int 值的 POJO:-

import java.io.*;

public class TestingSQLPerson implements Serializable {

// Here is Data in a comma Separated file: -
// Sumit,20,123455
// Ramit,40,12345

private String name;
private int age;
private Long testL;

public Long getTestL() {
    return testL;
}

public void setTestL(Long testL) {
    this.testL = testL;
}

public String getName() {
    return name;
}

public void setName(String name) {
    this.name = name;
}

public int getAge() {
    return age;
}

public void setAge(int age) {
    this.age = age;
}

}

这里是 Java 中的 Spark SQL 代码:-

import org.apache.spark.*;
import org.apache.spark.sql.*;
import org.apache.spark.api.java.*;
import org.apache.spark.api.java.function.Function;

public class TestingLongSQLTypes {

public static void main(String[] args) {

    SparkConf javaConf = new SparkConf();
    javaConf.setAppName("Test Long TTyypes");
    JavaSparkContext javaCtx = new JavaSparkContext(javaConf);


    SQLContext sqlContext = new org.apache.spark.sql.SQLContext(javaCtx);

    String dataFile = "file:///home/ec2-user/softwares/crime-data/testfile.txt";
    JavaRDD<TestingSQLPerson> people = javaCtx.textFile(dataFile).map(
      new Function<String, TestingSQLPerson>() {
        public TestingSQLPerson call(String line) throws Exception {
          String[] parts = line.split(",");

          TestingSQLPerson person = new TestingSQLPerson();
          person.setName(parts[0]);
          person.setAge(Integer.parseInt(parts[1].trim()));
          person.setTestL(Long.parseLong(parts[2].trim()));

          return person;
        }
      });

    // Apply a schema to an RDD of JavaBeans and register it as a table.
    DataFrame schemaPeople = sqlContext.createDataFrame(people, TestingSQLPerson.class);
    schemaPeople.registerTempTable("TestingSQLPerson");

    schemaPeople.printSchema();
    schemaPeople.show();


}

}

以上所有工作,最后在驱动程序控制台上,我可以看到没有任何异常错误的结果。 @Yukti - 在您的情况下,它也应该可以工作,只要您按照上述示例中定义的相同步骤进行操作。万一有什么偏差,那就告诉我,我可以帮你试试。

【讨论】:

  • 根据 Yukti 问题:“在流式上下文中,我得到 SQLContext 如下”,即 SQLContext 是在 DStream 的转换之一中创建的。我尝试了这种情况并面临同样的问题 @Sumit跨度>
  • 是的,我认为 SQLContext 从 @Sumit 和 Yukti 代码之间的关键区别的流上下文中获得。我也在做流媒体,需要获取 SQLContext 才能使用 createDataFrame。也遇到了同样的问题。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-07-18
  • 1970-01-01
  • 2017-04-27
  • 2018-12-17
相关资源
最近更新 更多