【问题标题】:Map reduce to perform group by and sum in Cassandra, with spark and job server映射减少以在 Cassandra 中执行分组和求和,带有火花和作业服务器
【发布时间】: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


【解决方案1】:

您的 map1 是一个 RDD。您可以尝试以下方法:

map1.foreach(r => println(r))

【讨论】:

    【解决方案2】:

    Spark 对 rdd 进行惰性求值。所以尝试一些操作

       map1.take(10).foreach(println)
    

    【讨论】:

      猜你喜欢
      • 2019-04-03
      • 2018-01-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-12-24
      • 2017-05-13
      • 1970-01-01
      • 2023-04-03
      相关资源
      最近更新 更多