如何在spark数据集上使用group by

Swa*_*shi 5 dataset apache-spark apache-spark-dataset

我正在使用Spark Dataset(Spark 1.6.1版本).以下是我的代码

object App { 

val conf = new SparkConf()
.setMaster("local")
.setAppName("SparkETL")

val sc = new SparkContext(conf)
sc.setLogLevel("ERROR")
val sqlContext = new SQLContext(sc);
import sqlContext.implicits._

}

override def readDataTable(tableName:String):DataFrame={
val dataFrame= App.sqlContext.read.jdbc(JDBC_URL, tableName, JDBC_PROP);
return dataFrame;
}


case class Student(stud_id , sname , saddress)
case class Student(classid, stud_id, name)


var tbl_student = JobSqlDAO.readDataTable("tbl_student").filter("stud_id = '" + studId + "'").as[Student].as("tbl_student")

var tbl_class_student = JobSqlDAO.readDataTable("tbl_class_student").as[StudentClass].as("tbl_class_student")


 var result = tbl_class_student.joinWith(tbl_student, $"tbl_student.stud_id" === $"tbl_class_student.stud_id").as("ff")
Run Code Online (Sandbox Code Playgroud)

现在我想在多列上执行group by子句?怎么做? result.groupBy(_._1._1.created_at)我可以这样做吗?如果是的话,那么我不能将结果看作一个组也是如何在多个列上进行的?

Vin*_*gio 0

如果我正确理解了您的要求,那么您最好的选择是使用PairRDDFunctionsreduceByKey类中的函数。

\n\n

该函数的签名是 def reduceByKey(func: (V, V) \xe2\x87\x92 V): RDD[(K, V)],它仅意味着您使用一系列键/值对。

\n\n

让我解释一下工作流程:

\n\n
    \n
  1. 您检索您想要使用的集合(在您的代码中result:)
  2. \n
  3. 使用 RDDmap函数,您可以将结果集拆分为一个元组,该元组包含两个子元组,其中包含组成键的字段和要聚合的字段(例如result.map(row => ((row.key1, row.key2), (row.value1, row.value2)):)
  4. \n
  5. 现在你有一个 RDD[(K,V)],其中类型 K 是键字段元组的类型,V 是值字段元组的类型
  6. \n
  7. 您可以通过传递聚合值reduceByKey类型的函数来直接使用(例如:)(V,V) => V(agg: (Int, Int), val: (Int, Int)) => (agg._1 + val._1, agg._2 + val._2)
  8. \n
\n\n

请注意:

\n\n
    \n
  • 您必须从聚合函数返回相同的值类型
  • \n
  • 您必须导入org.apache.spark.SparkContext._才能自动使用 PairRDDFunctions 实用函数
  • \n
  • 同样的推理也适用于groupBy,您必须从起始 RDD 映射到一对RDD[K,V],但您没有聚合函数,因为您只是将值存储在 seq 中以进行进一步计算
  • \n
  • 如果您需要聚合的起始值(例如:0 用于计数),请使用foldByKey函数
  • \n
\n