【发布时间】:2016-09-01 08:17:59
【问题描述】:
我正在创建一个连接到 cassandra 的 spark 作业服务器。获得记录后,我想执行一个简单的分组并对其求和。我能够检索数据,但无法打印输出。我已经在 google 上尝试了几个小时,并且也在 cassandra google 群组中发帖。我当前的代码如下,我在收集时遇到错误。
override def runJob(sc: SparkContext, config: Config): Any = {
//sc.cassandraTable("store", "transaction").select("terminalid","transdate","storeid","amountpaid").toArray().foreach (println)
// Printing of each record is successful
val rdd = sc.cassandraTable("POSDATA", "transaction").select("terminalid","transdate","storeid","amountpaid")
val map1 = rdd.map ( x => (x.getInt(0), x.getInt(1),x.getDate(2))->x.getDouble(3) ).reduceByKey((x,y)=>x+y)
println(map1)
// output is ShuffledRDD[3] at reduceByKey at Daily.scala:34
map1.collect
//map1.ccollectAsMap().map(println(_))
//Throwing error java.lang.ClassNotFoundException: transaction.Daily$$anonfun$2
}
【问题讨论】:
-
工作节点上是否有 spark cassandra 连接器运行时库?
-
请记住,Spark 是惰性的 - 在调用最终操作(如收集、获取、foreach 等)之前不会应用转换。因此,println 不会强制进行任何计算,它只是在 RDD 上调用 toString。因此,您无法确定是否已检索到该数据
-
@noorul 我有 cassandra 连接驱动程序。下面一行是打印记录“ sc.cassandraTable("store", "transaction").select("terminalid","transdate","storeid","amountpaid").toArray().foreach (println)"跨度>
标签: scala apache-spark cassandra spark-jobserver