☰
Flink Hive 方言中的 SORT BY / DISTRIBUTE BY / CLUSTER BY:分区排序语义与底层执行原理
2026/9/25 11:02:09 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

Flink 的 Hive 方言(Hive dialect)为了兼容 HiveQL 脚本迁移,提供了SORT BY、DISTRIBUTE BY与CLUSTER BY三类分区级排序子句,它们与ORDER BY的全局排序语义有本质区别。本文基于 Flink 仓库中的 Hive 方言文档 sort-cluster-distribute-by.md,完整讲解这三个子句的语法、参数与示例,并结合 LogicalDistribution、FlinkLogicalDistribution 与 BatchPhysicalDistributionRule 等规划器源码,剖析这些子句是如何被翻译成分区哈希分布加局部排序的物理算子组合的。读完本文,你能准确在 Hive 方言下编写分区排序查询,并理解其与ORDER BY的语义差异在物理计划中的落点。

一、背景:Hive 方言查询语法中的位置

Hive 方言支持 Hive DQL 的一个常用子集,完整的 SELECT 语法骨架定义在查询总览文档 overview.md 中,其中排序相关子句处于如下位置:

[WITH CommonTableExpression [ , ... ]] SELECT [ALL | DISTINCT] select_expr [ , ... ] FROM table_reference [WHERE where_condition] [GROUP BY col_list] [ORDER BY col_list] [CLUSTER BY col_list | [DISTRIBUTE BY col_list] [SORT BY col_list] ] [LIMIT [offset,] rows]

从语法骨架可以直接看出三者的关系:ORDER BY与CLUSTER BY互斥,而DISTRIBUTE BY可以与SORT BY组合使用,但CLUSTER BY与DISTRIBUTE BY/SORT BY的组合是二选一的关系。此外要注意一个重要前提:Hive 方言不再支持标准 Flink SQL 查询,若需写 Flink 语法应切换回默认方言(default dialect)。

二、SORT BY:仅保证分区内有序

语义描述

与 ORDER BY 保证输出的全局总序不同,SORT BY只保证每个分区内部的行按用户指定的顺序排列。因此当存在多个分区时,SORT BY返回的结果只是部分有序(partially ordered)的。

这一差异在分布式执行场景下的代价完全不同:ORDER BY要求最终由单一任务对全部输出排序(overview 文档中明确警告:当输出行数过大时可能耗费极长时间),而SORT BY的排序发生在各分区内部,天然可并行。

语法

query: SELECT expression [ , ... ] FROM src sortBy sortBy: SORT BY expression colOrder [ , ... ] colOrder: ( ASC | DESC )

参数

  • colOrder:指定返回行的排序方向,默认值为ASC。

示例

SELECT x, y FROM t SORT BY x; SELECT x, y FROM t SORT BY abs(y) DESC;

注意排序表达式可以是任意表达式(如abs(y)),而非仅限列引用。

三、DISTRIBUTE BY:重分区(repartition)

语义描述

DISTRIBUTE BY子句用于对数据进行重分区:指定表达式求值结果相同的行,会被划分到同一个分区中。它本身不产生任何排序保证,只控制数据在分区间的分布方式。

语法

distributeBy: DISTRIBUTE BY expression [ , ... ] query: SELECT expression [ , ... ] FROM src distributeBy

示例

-- 仅使用 DISTRIBUTE BY 子句 SELECT x, y FROM t DISTRIBUTE BY x; SELECT x, y FROM t DISTRIBUTE BY abs(y); -- 同时使用 DISTRIBUTE BY 和 SORT BY 子句 SELECT x, y FROM t DISTRIBUTE BY x SORT BY y DESC;

最后一条语句展示了组合用法:先按x求值结果将数据哈希分发到各分区,再在每个分区内部按y降序排序。

四、CLUSTER BY:DISTRIBUTE BY 与 SORT BY 的简写

语义描述

CLUSTER BY是DISTRIBUTE BY与SORT BY的组合简写:它先基于输入表达式对数据重分区,再在每个分区内对数据排序。同样地,该子句只保证数据在每个分区内有序,不保证全局顺序。

语法

clusterBy: CLUSTER BY expression [ , ... ] query: SELECT expression [ , ... ] FROM src clusterBy

示例

SELECT x, y FROM t CLUSTER BY x; SELECT x, y FROM t CLUSTER BY abs(y);

CLUSTER BY x等价于DISTRIBUTE BY x SORT BY x。

五、端到端示例:在 Hive 方言下执行 CLUSTER BY

overview 文档给出了一个可复制运行的完整会话,演示了从建立 Hive Catalog、加载 hive 模块、切换方言到执行CLUSTER BY的全过程:

Flink SQL> create catalog myhive with ('type' = 'hive', 'hive-conf-dir' = '/opt/hive-conf'); [INFO] Execute statement succeeded. Flink SQL> use catalog myhive; [INFO] Execute statement succeeded. Flink SQL> load module hive; [INFO] Execute statement succeeded. Flink SQL> use modules hive,core; [INFO] Execute statement succeeded. Flink SQL> set table.sql-dialect=hive; [INFO] Session property has been set. FLINK SQL> set sql-client.execution.result-mode=tableau; Flink SQL> select explode(array(1,2,3)); -- 调用 hive udtf +----+-------------+ || op | col | +----+-------------+ || +I | 1 | || +I | 2 | || +I | 3 | +----+-------------+ Received a total of 3 rows Flink SQL> create table tbl (key int,value string); [INFO] Execute statement succeeded. Flink SQL> insert into table tbl values (5,'e'),(1,'a'),(1,'a'),(3,'c'),(2,'b'),(3,'c'),(3,'c'),(4,'d'); [INFO] Submitting SQL update statement to the cluster... [INFO] SQL update statement has been successfully submitted to the cluster: FLINK SQL> set execution.runtime-mode=batch; -- 切换到批模式 Flink SQL> select * from tbl cluster by key; -- 执行 cluster by +-----+-------+ || key | value | +-----+-------+ || 1 | a | || 1 | a | || 5 | e | || 2 | b | || 3 | c | || 3 | c | || 3 | c | || 4 | d | +-----+-------+ Received a total of 8 rows

注意结果中key=1的行位于开头、key=5紧随其后,而2/3/4的行交错出现——这正是"分区内有序、分区间无序"语义的直观体现:同一key的行被哈希到同一分区后在该分区内有序,但各分区输出结果的拼接顺序并不保证全局升序。适用前提是切换到批模式(execution.runtime-mode=batch),这与下文源码分析中"物理转换规则位于 batch 规则集"相印证。

六、源码原理:从 SQL 子句到物理算子

6.1 逻辑计划节点 LogicalDistribution

在规划器中,这三个子句被统一建模为一个专门的逻辑节点。LogicalDistribution.java 的类注释直接说明了其定位:

/** * LogicalDistribution is used to represent the expected distribution of the data, similar to Hive's * SORT BY, DISTRIBUTE BY, and CLUSTER BY semantics. */ public class LogicalDistribution extends SingleRel { // distribution keys private final List<Integer> distKeys; // sort collation private final RelCollation collation;

该节点携带两个核心信息:

  • distKeys(分布键列索引列表):对应DISTRIBUTE BY/CLUSTER BY的表达式列,决定哈希分区的键;
  • collation(排序规格):对应SORT BY/CLUSTER BY的排序方向与列序。

由此可以推断三者的内部表达:DISTRIBUTE BY只填distKeys,SORT BY只填collation,CLUSTER BY则同时填充两者——与文档描述完全一致。

6.2 转换为 Flink 逻辑节点并派生分布特性

FlinkLogicalDistribution.scala 负责把上述通用逻辑节点转换为 Flink 规划器的逻辑节点,并在create方法中根据distKeys是否为空派生出不同的分布特性(distribution trait):

val traitSet = if (distKeys.isEmpty) { cluster .traitSetOf(FlinkConventions.LOGICAL) .replace(collationTrait) .replace(FlinkRelDistribution.ANY) } else { cluster .traitSetOf(FlinkConventions.LOGICAL) .replace(collationTrait) .replace(FlinkRelDistribution.hash(distKeys)) }

这段代码精确对应了文档语义:

  • 无分布键(仅SORT BY):分布特性为ANY,即不强制重分区,只需在任意现有分区内满足排序要求——对应"每个分区内有序"的文档描述;
  • 有分布键(DISTRIBUTE BY或CLUSTER BY):分布特性为hash(distKeys),即按分布键哈希重分区——对应"相同表达式求值结果的行进入同一分区"。

6.3 物理转换规则:哈希 Exchange + 局部 Sort

BatchPhysicalDistributionRule.scala 将FlinkLogicalDistribution转换为批物理算子,其核心逻辑为:

val requiredTraitSet = input.getTraitSet .replace(distribution) // 要求输入满足目标分布(hash 或 ANY) .replace(FlinkConventions.BATCH_PHYSICAL) val newInput = RelOptRule.convert(input, requiredTraitSet) if (logicalDistribution.collation.getFieldCollations.isEmpty) { newInput // 无排序要求:仅需重分区(或原样) } else { new BatchPhysicalSort( // 有排序要求:在满足分布的输入上做局部排序 logicalDistribution.getCluster, providedTraitSet, newInput, logicalDistribution.collation) }

从源码结构看,物理执行计划由此生成:

  1. 规划器先通过RelOptRule.convert让子计划满足要求的分布特性——当分布特性是hash(distKeys)时,这会在计划中引入一次按哈希键重分区的 Exchange;当是ANY时则不引入额外重分区;
  2. 随后仅当collation非空(即存在SORT BY或CLUSTER BY)时,才在满足分布的输入之上追加BatchPhysicalSort,且排序作用域限定在重分区后的各分区内部。

这条转换链解释了文档中"CLUSTER BY 先重分区、再分区内排序"的执行流程,也解释了为何它不产生全局序:BatchPhysicalSort的排序发生在上游哈希分区之后,各分区独立排序,分区之间不再汇合排序(对比之下ORDER BY才需要汇合到单任务做全局排序,正如 overview 文档的 warning 所述)。

七、实践要点小结

  1. 选择SORT BY:只想在现有分区内排序、不关心数据重新分布时最便宜,对应FlinkRelDistribution.ANY分布特性,不引入额外 Exchange;
  2. 选择DISTRIBUTE BY:需要让相同键的行落入同一分区(例如为下游按分区处理做准备)但不需要排序时;
  3. 选择CLUSTER BY:同时需要"按键重分区 + 分区内排序"时,它是DISTRIBUTE BY ... SORT BY ...(按键同时排序)的简写;
  4. 需要全局有序结果时使用ORDER BY:三者均不保证跨分区的全局顺序,而ORDER BY以单任务全局排序为代价提供该保证,数据量大时应慎用;
  5. 运行前提:这些子句属于 Hive 方言(table.sql-dialect=hive)下的批查询能力,物理转换规则位于批规则集中,建议配合execution.runtime-mode=batch使用;编写 Flink 原生语法时应切回默认方言。

以上行为均以当前仓库中的文档与规划器源码为准:语义描述来自 sort-cluster-distribute-by.md 与 overview.md,执行原理证据来自 LogicalDistribution.java、FlinkLogicalDistribution.scala 与 BatchPhysicalDistributionRule.scala。

  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

上一篇:OpenPLC Editor:5个理由让你立即上手的开源PLC编程平台
下一篇:3分钟上手!免费开源PLC编程软件OpenPLC Editor完全指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询