大数据流处理批处理数据工程【免费下载链接】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 BYDISTRIBUTE 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 BYoverview 文档给出了一个可复制运行的完整会话演示了从建立 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-dialecthive; [INFO] Session property has been set. FLINK SQL set sql-client.execution.result-modetableau; 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-modebatch; -- 切换到批模式 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注意结果中key1的行位于开头、key5紧随其后而2/3/4的行交错出现——这正是分区内有序、分区间无序语义的直观体现同一key的行被哈希到同一分区后在该分区内有序但各分区输出结果的拼接顺序并不保证全局升序。适用前提是切换到批模式execution.runtime-modebatch这与下文源码分析中物理转换规则位于 batch 规则集相印证。六、源码原理从 SQL 子句到物理算子6.1 逻辑计划节点 LogicalDistribution在规划器中这三个子句被统一建模为一个专门的逻辑节点。LogicalDistribution.java 的类注释直接说明了其定位/** * LogicalDistribution is used to represent the expected distribution of the data, similar to Hives * SORT BY, DISTRIBUTE BY, and CLUSTER BY semantics. */ public class LogicalDistribution extends SingleRel { // distribution keys private final ListInteger distKeys; // sort collation private final RelCollation collation;该节点携带两个核心信息distKeys分布键列索引列表对应DISTRIBUTE BY/CLUSTER BY的表达式列决定哈希分区的键collation排序规格对应SORT BY/CLUSTER BY的排序方向与列序。由此可以推断三者的内部表达DISTRIBUTE BY只填distKeysSORT BY只填collationCLUSTER BY则同时填充两者——与文档描述完全一致。6.2 转换为 Flink 逻辑节点并派生分布特性FlinkLogicalDistribution.scala 负责把上述通用逻辑节点转换为 Flink 规划器的逻辑节点并在create方法中根据distKeys是否为空派生出不同的分布特性distribution traitval 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 局部 SortBatchPhysicalDistributionRule.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) }从源码结构看物理执行计划由此生成规划器先通过RelOptRule.convert让子计划满足要求的分布特性——当分布特性是hash(distKeys)时这会在计划中引入一次按哈希键重分区的 Exchange当是ANY时则不引入额外重分区随后仅当collation非空即存在SORT BY或CLUSTER BY时才在满足分布的输入之上追加BatchPhysicalSort且排序作用域限定在重分区后的各分区内部。这条转换链解释了文档中CLUSTER BY 先重分区、再分区内排序的执行流程也解释了为何它不产生全局序BatchPhysicalSort的排序发生在上游哈希分区之后各分区独立排序分区之间不再汇合排序对比之下ORDER BY才需要汇合到单任务做全局排序正如 overview 文档的 warning 所述。七、实践要点小结选择SORT BY只想在现有分区内排序、不关心数据重新分布时最便宜对应FlinkRelDistribution.ANY分布特性不引入额外 Exchange选择DISTRIBUTE BY需要让相同键的行落入同一分区例如为下游按分区处理做准备但不需要排序时选择CLUSTER BY同时需要按键重分区 分区内排序时它是DISTRIBUTE BY ... SORT BY ...按键同时排序的简写需要全局有序结果时使用ORDER BY三者均不保证跨分区的全局顺序而ORDER BY以单任务全局排序为代价提供该保证数据量大时应慎用运行前提这些子句属于 Hive 方言table.sql-dialecthive下的批查询能力物理转换规则位于批规则集中建议配合execution.runtime-modebatch使用编写 Flink 原生语法时应切回默认方言。以上行为均以当前仓库中的文档与规划器源码为准语义描述来自 sort-cluster-distribute-by.md 与 overview.md执行原理证据来自 LogicalDistribution.java、FlinkLogicalDistribution.scala 与 BatchPhysicalDistributionRule.scala。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink Hive 方言详解SORT BY、DISTRIBUTE BY 与 CLUSTER BY 的语义、语法与源码实现Flink Hive 方言详解SORT BY、DISTRIBUTE BY 与 CLUSTER BY 的语义、语法与源码实现 本篇聚焦 Flink 表模块 Hi大数据流处理批处理数据工程Apache Spark SQL 的 SORT BY 子句分区内排序语法、NULL 排序语义与底层实现解析Apache Spark SQL 的 SORT BY 子句分区内排序语法、NULL 排序语义与底层实现解析 输出文章 标签内的这篇技术指南完整讲解 Apac大数据数据分析批处理流处理机器学习图计算Apache Spark SQL DISTRIBUTE BY 子句详解按表达式重分区与 CLUSTER BY 的区别Apache Spark SQL DISTRIBUTE BY 子句详解按表达式重分区与 CLUSTER BY 的区别 导读 本文是 Apache Spark大数据数据分析批处理流处理机器学习图计算上一篇OpenPLC Editor5个理由让你立即上手的开源PLC编程平台下一篇3分钟上手免费开源PLC编程软件OpenPLC Editor完全指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考