0Pricing
Scala for Backend Engineering & Functional Programming · 课时

聚合

分组并聚合

聚合 是 CoddyKit 上的免费 Scala for Backend Engineering & Functional Programming 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Scala for Backend Engineering & Functional Programming 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。

为什么要聚合?

聚合会将许多行概括为较少的结果:总计、平均值以及每个分组的计数。在 Spark 中,聚合属于宽操作,可能会在集群中跨节点洗牌数据。

全局聚合

使用 agg 以及 count、sum、avg、min、max 等函数,对整个 DataFrame 计算单个汇总结果。

import org.apache.spark.sql.functions._

df.agg(
  count("*").as("rows"),
  avg("age").as("avg_age")
).show()

groupBy

groupBy 会根据一个或多个列对行进行分区,并生成可用于聚合的 RelationalGroupedDataset。

import org.apache.spark.sql.functions._

val byCity = df.groupBy("city").agg(count("*").as("people"))
byCity.show()

多个聚合

向 agg 传入多个聚合表达式,即可一次遍历为每个分组计算多个汇总结果。

import org.apache.spark.sql.functions._

df.groupBy("department").agg(
  sum("salary").as("total"),
  avg("salary").as("avg"),
  max("salary").as("top")
).show()

统计不同值

使用 countDistinct 统计唯一值;对于超大数据集,则可以使用 approx_count_distinct 更快地进行近似计数。

import org.apache.spark.sql.functions._

df.agg(countDistinct("city").as("unique_cities")).show()

使用 having 筛选分组

在 SQL 中,HAVING 会筛选聚合后的分组。在 DSL 中,请在 agg 之后对计算出的列应用 filter。

import org.apache.spark.sql.functions._

df.groupBy("city")
  .agg(count("*").as("n"))
  .filter(col("n") > 100)
  .show()

在 SQL 中进行聚合

对已注册的视图执行 SQL 查询,实现相同的逻辑,并使用 GROUP BY 和 HAVING。

df.createOrReplaceTempView("people")
spark.sql(
  "SELECT city, COUNT(*) AS n FROM people GROUP BY city HAVING COUNT(*) > 100"
).show()

透视表

pivot 会将某列的不同值转换为单独的列,适合制作交叉表,例如按地区和季度统计销售额。

import org.apache.spark.sql.functions._

sales.groupBy("region")
     .pivot("quarter")
     .agg(sum("amount"))
     .show()

窗口函数

窗口函数会在滑动窗口上进行聚合,而不会合并行,非常适合计算累计总额或排名。

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

val w = Window.partitionBy("dept").orderBy(col("salary").desc)
df.withColumn("rank", rank().over(w)).show()

RDD 聚合

在 RDD 层级,aggregateByKey 和 reduceByKey 会使用自定义逻辑按键合并值,并通过本地预合并提高效率。

val pairs = sc.parallelize(Seq(("a", 10), ("a", 20), ("b", 5)))
val sums  = pairs.reduceByKey(_ + _)
// (a, 30), (b, 5)

纯 Scala groupBy

这是一个自包含的类比示例:Scala 集合的 groupBy 加上 mapValues,对应 Spark 中的分组与聚合。

object Main {
  def main(args: Array[String]): Unit = {
    val data = Seq(("a", 10), ("a", 20), ("b", 5))
    val sums = data.groupBy(_._1).map { case (k, v) => k -> v.map(_._2).sum }
    sums.toSeq.sortBy(_._1).foreach { case (k, s) => println(s"$k: $s") }
  }
}

快速检查

哪个功能可以在一个行窗口上进行聚合,而不会将这些行合并为每组一行?

回顾

您在 Spark 中完成了数据聚合:

  • 使用 agg 配合 count、sum、avg、min、max
  • 使用 groupBy、多重聚合和 having 筛选
  • 使用 pivot 和窗口函数
  • 使用 RDD 的 reduceByKey / aggregateByKey

您已完成 Apache Spark 课程。

常见问题解答

「聚合」课时是免费的吗?

是的 — 「聚合」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Scala for Backend Engineering & Functional Programming 课程的其余内容,请升级到 CoddyKit PRO。 Scala for Backend Engineering & Functional Programming 课程共包含 4 节课。

「聚合」这节课中我会学到什么?

分组并聚合 你通过在浏览器中直接运行的动手代码来练习 Scala for Backend Engineering & Functional Programming,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Scala for Backend Engineering & Functional Programming 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Scala for Backend Engineering & Functional Programming 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。

「聚合」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Scala for Backend Engineering & Functional Programming 课中编写并运行代码吗?

能。每节 Scala for Backend Engineering & Functional Programming 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. RDD 与 DataFrames
  2. 转换与操作
  3. Spark SQL
  4. 聚合
← 返回 Scala for Backend Engineering & Functional Programming