Apache Spark SQL CLUSTER BY 子句完全指南:分区内聚类与排序的语义、实现与实战
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
CLUSTER BY是 Apache Spark SQL 中用于按表达式重新分区数据,并在每个分区内部完成排序的查询组织子句,其语义等价于DISTRIBUTE BY加上SORT BY。本文基于当前仓库的官方文档 docs/sql-ref-syntax-qry-select-clusterby.md 与 Catalyst 解析器、逻辑计划源码,系统讲解其语法、参数、与ORDER BY/SORT BY/DISTRIBUTE BY的异同、底层实现原理以及可复制的完整示例,帮助你正确使用该子句完成数据倾斜预处理、ETL 落盘优化等场景中的分区聚类与局部排序。
概述:CLUSTER BY 做了什么
CLUSTER BY子句执行两步操作:
- 重新分区(repartition):根据输入表达式对数据进行重分布,使得具有相同表达式值(或其哈希值相同)的行被放入同一个分区;
- 分区内排序(sort):在重分区完成后,对每个分区内部的数据按照这些表达式进行排序。
因此CLUSTER BY在语义上完全等价于先执行 DISTRIBUTE BY(负责重分区)再执行 SORT BY(负责分区内排序)。
这里有一个关键前提需要牢记:CLUSTER BY只保证结果行在每个分区内部是有序的,并不保证整个输出是全局有序的。如果需要全局有序的输出,必须使用 ORDER BY 子句。
语法
CLUSTER BY { expression [ , ... ] }- 支持指定一个或多个表达式;
- 多个表达式之间以逗号分隔;
CLUSTER BY出现在SELECT语句的查询组织(query organization)部分,通常位于WHERE、GROUP BY、HAVING等子句之后。
参数说明:expression
| 参数 | 说明 |
|---|---|
expression | 指定一个或多个值、运算符和 SQL 函数组合而成的计算结果。CLUSTER BY既按该表达式的结果对数据进行重分区,也按该表达式的值在分区内排序。 |
表达式可以是简单列名(如age),也可以是任意可求值的表达式,例如算术运算、字符串函数、CASE表达式等。在实际查询中,它既可以引用SELECT列表中出现的列,也可以引用源表中未出现在结果集中的列。
完整示例:从无排序到分区聚类
下面的示例完整取自官方文档,通过把 shuffle 分区数降到 2,可以更直观地观察CLUSTER BY的聚类与排序行为。
CREATE TABLE person (name STRING, age INT); INSERT INTO person VALUES ('Zen Hui', 25), ('Anil B', 18), ('Shone S', 16), ('Mike A', 25), ('John A', 18), ('Jack N', 16); -- 将 shuffle 分区数降为 2,便于观察 CLUSTER BY 的聚类与排序行为 SET spark.sql.shuffle.partitions = 2;第一步:不加任何排序子句的普通查询。没有任何排序指令时,查询结果是不确定的,age列并未排序:
SELECT age, name FROM person; +---+-------+ |age| name| +---+-------+ | 16|Shone S| | 25|Zen Hui| | 16| Jack N| | 25| Mike A| | 18| John A| | 18| Anil B| +---+-------+第二步:使用CLUSTER BY age。相同age的人员被聚到同一分区,且每个分区内部按age升序排列。在示例的输出中,年龄为 18 和 25 的人员位于第一个分区,年龄为 16 的人员位于第二个分区:
SELECT age, name FROM person CLUSTER BY age; +---+-------+ |age| name| +---+-------+ | 18| John A| | 18| Anil B| | 25|Zen Hui| | 25| Mike A| | 16|Shone S| | 16| Jack N| +---+-------+对比两个结果可以发现:使用CLUSTER BY后,相同age值的行必然相邻(聚类),且每个分区内的行按age有序;但整个结果集依然是分区间的"部分有序"——第一个分区输出18/25后,第二个分区才输出16,全局来看并非严格升序。
等价写法验证:DISTRIBUTE BY + SORT BY
由于CLUSTER BY age语义上等价于先DISTRIBUTE BY age再SORT BY age,你完全可以用下面的写法得到等价的结果:
SELECT age, name FROM person DISTRIBUTE BY age SORT BY age;值得注意的区别是:DISTRIBUTE BY只做重分区、不排序,其示例输出中分区内行的顺序是随机的;SORT BY只做分区内排序、不改变分区方式;CLUSTER BY则是两者行为的叠加,并且分区内排序方向固定为升序。
源码视角:CLUSTER BY 在 Catalyst 中如何实现
从当前仓库的源码可以确认CLUSTER BY的完整执行路径,这有助于理解其"重分区 + 局部排序"的底层语义。
解析阶段:AstBuilder 生成逻辑计划
在 sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/parser/AstBuilder.scala 的withQueryResultClauses方法中,CLUSTER BY被翻译为如下逻辑计划:
} else if (order.isEmpty && sort.isEmpty && distributeBy.isEmpty && !clusterBy.isEmpty) { clause = PipeOperators.clusterByClause val expressions = expressionList(clusterBy) Sort( expressions.map(SortOrder(_, Ascending)), global = false, withRepartitionByExpression(ctx, expressions, query)) }这段代码清晰地揭示了三点实现事实:
- 重分区算子:通过
withRepartitionByExpression生成RepartitionByExpression逻辑算子; - 分区内排序:通过
Sort(..., global = false, ...)生成非全局排序(global = false正是"仅在分区内排序"的源码级体现); - 排序方向固定:每个表达式被包装成
SortOrder(_, Ascending),即CLUSTER BY的排序方向固定为升序,这一点与SORT BY(可显式指定ASC/DESC以及NULLS FIRST/NULLS LAST)不同。
同时,该解析方法还规定ORDER BY、SORT BY、DISTRIBUTE BY、CLUSTER BY是互斥的——如果在一个查询中同时出现多个这类子句,会抛出combinationQueryResultClausesUnsupportedError。
逻辑计划阶段:RepartitionByExpression 算子
重分区算子定义在 sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala:
case class RepartitionByExpression( partitionExpressions: Seq[Expression], child: LogicalPlan, optNumPartitions: Option[Int], optAdvisoryPartitionSize: Option[Long] = None) extends RepartitionOperation with HasPartitionExpressions { override def shuffle: Boolean = true ... }关键点是override def shuffle: Boolean = true:CLUSTER BY必然触发一次shuffle(即数据重分布),这意味着它涉及网络传输与磁盘读写开销,与无需 shuffle 的ORDER BY(单分区场景除外)在代价模型上并不相同。对于大规模数据,shuffle 的开销需要通过调整spark.sql.shuffle.partitions等参数来控制并行度与分区大小。
物理执行角度
RepartitionByExpression在物理执行阶段会转换为基于表达式哈希值(HashPartitioning)的Exchange算子;数据到达各分区后,再执行局部的排序。因此可以推断:具有相同表达式值的行必然落在同一分区,但多个分区之间不存在全局顺序约束,这正是文档所述"只保证分区内有序、不保证全局有序"的实现根源。
CLUSTER BY 与 ORDER BY / SORT BY / DISTRIBUTE BY 的对照
| 子句 | 是否重分区(shuffle) | 是否分区内排序 | 是否全局排序 | 排序方向控制 |
|---|---|---|---|---|
CLUSTER BY | 是 | 是 | 否 | 固定升序 |
DISTRIBUTE BY | 是 | 否 | 否 | 不适用 |
SORT BY | 否 | 是 | 否 | 可指定ASC/DESC、NULLS FIRST/LAST |
ORDER BY | 否(单分区可无 shuffle) | 是 | 是 | 可指定ASC/DESC、NULLS FIRST/LAST |
使用建议:
- 需要全局有序结果(如对外输出必须严格升序)→ 使用
ORDER BY; - 只想做分区内排序、保持现有分区→ 使用
SORT BY; - 只想把相同键的行聚合到同一分区、不关心顺序→ 使用
DISTRIBUTE BY; - 既要按键分区聚类、又要每个分区内有序→ 使用
CLUSTER BY,例如将数据按日期/用户 ID 聚类后再分区写入下游系统,可减少下游的局部排序与随机 IO。
进阶:管道操作符中的 CLUSTER BY
当前仓库已支持 SQL 管道操作符(pipe operators)语法,CLUSTER BY也可以作为管道操作符使用。在 sql/core/src/test/scala/org/apache/spark/sql/execution/SparkSqlParserSuite.scala 中存在如下测试:
checkRepartition("TABLE t |> CLUSTER BY x |> TABLESAMPLE (100 PERCENT)")这表明你可以写出这样的管道式查询:
TABLE person |> CLUSTER BY age;管道操作符语法中|> CLUSTER BY与普通SELECT ... CLUSTER BY ...在语义上一致,为流式组合多个查询步骤提供了更直观的写法。
相关子句速览
CLUSTER BY属于SELECT查询组织子句家族,与之紧密相关的文档包括:
- SELECT 主语句
- WHERE Clause
- GROUP BY Clause
- HAVING Clause
- ORDER BY Clause
- SORT BY Clause
- DISTRIBUTE BY Clause
- LIMIT Clause
- OFFSET Clause
- CASE Clause
- PIVOT Clause
- UNPIVOT Clause
- LATERAL VIEW Clause
小结
CLUSTER BY是"重分区聚类 + 分区内升序排序"二合一的高性价比子句,在数据倾斜预处理、ETL 分区落盘、减少下游局部排序等场景中非常实用。理解其三个要点即可正确使用:语义上等于DISTRIBUTE BY+SORT BY;只保证分区内有序、不保证全局有序;必然触发一次 shuffle 且排序方向固定为升序。
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考