【发布时间】:2015-03-23 04:51:11
【问题描述】:
我想根据我在 RDD 中的值从 Cassandra 查询一些数据。我的方法如下:
val userIds = sc.textFile("/tmp/user_ids").keyBy( e => e )
val t = sc.cassandraTable("keyspace", "users").select("userid", "user_name")
val userNames = userIds.flatMap { userId =>
t.where("userid = ?", userId).take(1)
}
userNames.take(1)
虽然 Cassandra 查询在 Spark shell 中工作,但当我在 flatMap 中使用它时会引发异常:
org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 2.0 failed 1 times, most recent failure: Lost task 0.0 in stage 2.0 (TID 2, localhost): java.lang.NullPointerException:
org.apache.spark.rdd.RDD.<init>(RDD.scala:125)
com.datastax.spark.connector.rdd.CassandraRDD.<init>(CassandraRDD.scala:49)
com.datastax.spark.connector.rdd.CassandraRDD.copy(CassandraRDD.scala:83)
com.datastax.spark.connector.rdd.CassandraRDD.where(CassandraRDD.scala:94)
我的理解是我无法在另一个 RDD 中生成 RDD(Cassandra 结果)。
我在网上找到的示例读取 RDD 中的整个 Cassandra 表并加入 RDD (像这样:https://cassandrastuff.wordpress.com/2014/07/07/cassandra-and-spark-table-joins/)。但如果 Cassandra 表很大,它就无法扩展。
但是我该如何解决这个问题呢?
【问题讨论】:
标签: scala cassandra apache-spark rdd