【发布时间】: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