Flink Hive 方言查询(Queries)完全指南:从 SELECT 语法到 Sort/Cluster/Join/CTE 实战

📅 发布时间:2026/9/24 8:42:06
Flink Hive 方言查询(Queries)完全指南:从 SELECT 语法到 Sort/Cluster/Join/CTE 实战
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读Hive 方言是 Flink 为兼容 Hive 生态提供的 SQL 解析与执行模式启用后你可以直接在 Flink 中编写 HiveQL 语法执行查询而无需在 Flink 与 Hive 之间来回切换。本文以 Hive 方言查询概览文档 为主体系统梳理 Flink Hive 方言支持的 DQL 子集——包括完整的SELECT语法骨架、WHERE/GROUP BY/ORDER BY/LIMIT等核心子句以及Sort/Cluster/Distribute By、Join、集合操作、Lateral View、窗口函数、子查询、CTE、Transform与Table Sample等扩展能力并结合仓库源码与集成测试说明其实现与使用前提。读完本文你将掌握在 Flink SQL Client / SQL Gateway / Table API 中开启 Hive 方言并编写各类查询的完整方法。说明Hive 方言的查询能力文档目录位于 docs/content.zh/docs/dev/table/hive-compatibility/hive-dialect/queries/本文是其总览Overview的深入展开各子句的详细语法见对应小节。Hive 方言查询能力概览Hive 方言支持 Hive DQL数据查询语言中常用子集覆盖以下查询特性各子句的完整文档链接见后文对应小节Sort/Cluster/Distributed BYGroup ByJoinSet OperationUNION / INTERSECT / EXCEPTLateral ViewWindow FunctionsSub-QueriesCTE公共表表达式TransformTable Sample使用前提先启用 Hive 方言在编写 Hive 方言查询之前需要先完成方言切换与相关环境准备。Flink 支持default与hive两种 SQL 方言你可以为执行的每条语句动态切换无需重启会话。以下关键前提来自 Hive 方言总览添加 Hive 相关依赖使用 Hive 方言前必须引入 Hive 依赖参考 Hive Dependencies。确保当前 Catalog 是 HiveCatalog否则将回退到 Flink 默认方言在启动了 HiveServer2 Endpoint 的 SQL Gateway 下默认 Catalog 即为 HiveCatalog。建议加载 HiveModule 并置于 Module 列表首位使函数解析优先命中 Hive 内置函数。标识符限制Hive 方言只支持db.table两级标识符不支持带 Catalog 名的标识符。运行模式限制Hive 方言主要在批Batch模式下使用Sort/Cluster/Distributed BY、Transform等语法尚未在流Streaming模式下支持。部分特性是否可用取决于 Hive 版本见 支持的 Hive 版本例如更新数据库位置仅在 Hive-2.4.0 或更高版本支持。SQL Client 中的切换方式通过table.sql-dialect属性指定方言Flink SQL SET table.sql-dialect hive; -- 切换到 Hive 方言 [INFO] Session property has been set. Flink SQL SET table.sql-dialect default; -- 切回 Flink 默认方言 [INFO] Session property has been set.SQL GatewayHiveServer2 Endpoint与 Table APISQL Gateway在配置了 HiveServer2 Endpoint 的 SQL Gateway 下HiveModule 已默认加载、默认 Catalog 即 HiveCatalog可通过SET table.sql-dialect hive切换方言。Table API可以通过EnvironmentSettings或执行配置指定方言同样支持逐语句切换。整体查询语法SELECT语句可以是包含 CTE、集合操作及其它各种子句的查询的一部分。Hive 方言的整体查询语法如下[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]语法要点SELECT语句可以是 集合查询 的一部分也可以是其它查询的 子查询CommonTableExpression是在WITH子句中指定查询派生的临时结果集详见 CTE 小节table_reference表示查询的输入可以是普通表、视图、Join 或 子查询表名和列名大小写不敏感。WHERE 子句WHERE条件是一个布尔表达式。Hive 方言在WHERE子句中支持 Hive 提供的大量操作符与 UDF并支持部分类型的子查询如IN/NOT IN/EXISTS/NOT EXISTS。GROUP BY 子句GROUP BY用于对多行输入结合给定聚合函数计算单个结果。详细语法含GROUPING SETS/ROLLUP/CUBE增强聚合见 GROUP BY 详解。ORDER BY 子句ORDER BY用于按用户指定顺序返回结果行。与 SORT BY 不同ORDER BY保证输出的全局有序。注意为保证全局有序最终排序必须由单个 task完成。因此如果输出行数过大可能耗时极长需谨慎使用。CLUSTER / DISTRIBUTE / SORT BY这三种子句只保证分区内有序/数据重分布详细说明见 Sort/Cluster/Distribute By 详解。ALL 与 DISTINCT 子句ALL与DISTINCT选项指定是否返回重复行两者都未指定时默认为ALL返回所有匹配行DISTINCT指定从结果集中去除重复行。LIMIT 子句LIMIT用于约束SELECT语句返回的行数接受一个或两个数值参数两者必须是非负整数常量第一个参数指定返回首行的偏移量offset第二个参数指定返回行的最大数量只给一个参数时它表示最大行数偏移量默认为 0。例如LIMIT 5返回前 5 行LIMIT 10, 5跳过前 10 行后返回 5 行。GROUP BY 详解Group by子句用于对给定的聚合函数从多行输入计算单个结果。Hive 方言还支持基于同一条记录执行多次聚合的增强聚合特性ROLLUP/CUBE/GROUPING SETS。语法group_by_clause: group_by_clause_1 | group_by_clause_2 group_by_clause_1: GROUP BY group_expression [ , ... ] [ WITH ROLLUP | WITH CUBE ] group_by_clause_2: GROUP BY { group_expression | { ROLLUP | CUBE | GROUPING SETS } ( grouping_set [ , ... ] ) } [ , ... ] grouping_set: { expression | ( [ expression [ , ... ] ] ) } groupByQuery: SELECT expression [ , ... ] FROM src groupByClause?在group_expression中列也可以按位置编号指定。但需注意 Hive 版本差异对于 Hive 0.11.0 至 2.1.x需设置hive.groupby.orderby.position.alias为 true默认 false对于 Hive 2.2.0 及以后需设置hive.groupby.position.alias为 true默认 false。GROUPING SETSGROUPING SETS允许比标准GROUP BY更复杂的分组操作行按每个指定分组集合分别分组聚合按组计算与简单的GROUP BY一致。所有GROUPING SETS子句都可以逻辑上等价地用多个由UNION连接的GROUP BY查询表达例如SELECT a, b, SUM( c ) FROM tab1 GROUP BY a, b GROUPING SETS ( (a, b), a, b, ( ) )等价于SELECT a, b, SUM( c ) FROM tab1 GROUP BY a, b UNION SELECT a, null, SUM( c ) FROM tab1 GROUP BY a, null UNION SELECT null, b, SUM( c ) FROM tab1 GROUP BY null, b UNION SELECT null, null, SUM( c ) FROM tab1当对某列显示聚合结果时其值为 null这可能与列本身包含 null 值冲突。此时GROUPING__ID函数是区分手段它返回一个位向量指示每一列是否存在——结果集中某行若该列已被聚合则产生 1否则为 0可用于区分数据中的 null。此外GROUPING函数指示GROUP BY子句中的表达式在给定行中是否被聚合值 0 表示该列属于分组集合值 1 表示该列不属于分组集合。ROLLUPROLLUP是常见分组集合类型的简写表示给定表达式列表及该列表的所有前缀包括空列表。例如GROUP BY a, b, c WITH ROLLUP等价于GROUP BY a, b, c GROUPING SETS ( (a, b, c), (a, b), (a), ( ) )CUBECUBE同样是常见分组集合类型的简写表示给定列表及其所有可能子集——即幂集。例如GROUP BY a, b, c WITH CUBE等价于GROUP BY a, b, c GROUPING SETS ( (a, b, c), (a, b), (b, c), (a, c), (a), (b), (c), ( ) )示例-- 按表达式分组 SELECT abs(x), sum(y) FROM t GROUP BY abs(x); -- 按列分组 SELECT x, sum(y) FROM t GROUP BY x; -- 按位置分组Hive 2.2.0 需开启 hive.groupby.position.alias SELECT x, sum(y) FROM t GROUP BY 1; -- 按表的第一列分组 -- 使用 grouping sets SELECT x, SUM(y) FROM t GROUP BY x GROUPING SETS ( x, ( ) ); -- 使用 rollup SELECT x, SUM(y) FROM t GROUP BY x WITH ROLLUP; SELECT x, SUM(y) FROM t GROUP BY ROLLUP (x); -- 使用 cube SELECT x, SUM(y) FROM t GROUP BY x WITH CUBE; SELECT x, SUM(y) FROM t GROUP BY CUBE (x);Sort/Cluster/Distribute By 详解SORT BY与保证输出全局有序的ORDER BY不同SORT BY只保证每个分区内的结果行按用户指定顺序排列。因此当存在多个分区时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;DISTRIBUTE BYDISTRIBUTE 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;CLUSTER BYCLUSTER BY是DISTRIBUTE BY与SORT BY的快捷组合先基于输入表达式重分区数据再对每个分区内排序。同样该子句只保证数据在每个分区内有序。clusterBy: CLUSTER BY expression [ , ... ] query: SELECT expression [ , ... ] FROM src clusterBySELECT x, y FROM t CLUSTER BY x; SELECT x, y FROM t CLUSTER BY abs(y);JOIN 详解JOIN用于基于连接条件合并两个关系的行。Hive 方言支持以下连接语法join_table: table_reference [ INNER ] JOIN table_factor [ join_condition ] | table_reference { LEFT | RIGHT | FULL } [ OUTER ] JOIN table_reference join_condition | table_reference LEFT SEMI JOIN table_reference [ ON expression ] | table_reference CROSS JOIN table_reference [ join_condition ] table_reference: table_factor | join_table table_factor: tbl_name [ alias ] | table_subquery alias | ( table_references ) join_condition: { ON expression | USING ( colName [, ...] ) }连接类型连接类型语义INNER JOIN返回两侧都匹配的行是默认连接类型LEFT JOIN返回左表全部行与右表匹配值无匹配则追加NULL等价于LEFT OUTER JOINRIGHT JOIN返回右表全部行与左表匹配值无匹配则追加NULL等价于RIGHT OUTER JOINFULL JOIN返回两侧全部行某一侧无匹配则追加NULL等价于FULL OUTER JOINLEFT SEMI JOIN只返回左表在右表中有匹配的行不拼接右表的值CROSS JOIN返回两侧的笛卡尔积示例-- INNER JOIN SELECT t1.x FROM t1 INNER JOIN t2 USING (x); SELECT t1.x FROM t1 INNER JOIN t2 ON t1.x t2.x; -- LEFT JOIN SELECT t1.x FROM t1 LEFT JOIN t2 USING (x); SELECT t1.x FROM t1 LEFT OUTER JOIN t2 ON t1.x t2.x; -- RIGHT JOIN SELECT t1.x FROM t1 RIGHT JOIN t2 USING (x); SELECT t1.x FROM t1 RIGHT OUTER JOIN t2 ON t1.x t2.x; -- FULL JOIN SELECT t1.x FROM t1 FULL JOIN t2 USING (x); SELECT t1.x FROM t1 FULL OUTER JOIN t2 ON t1.x t2.x; -- LEFT SEMI JOIN SELECT t1.x FROM t1 LEFT SEMI JOIN t2 ON t1.x t2.x; -- CROSS JOIN SELECT t1.x FROM t1 CROSS JOIN t2 USING (x);Set Operations集合操作详解集合操作用于将多个SELECT语句合并为单个结果集。Hive 方言支持UNION、INTERSECT、EXCEPT/MINUS三种操作。UNIONUNION/UNION DISTINCT/UNION ALL返回两侧都能找到的行UNION与UNION DISTINCT只返回去重后的行UNION ALL不去重。query { UNION [ ALL | DISTINCT ] } query [ .. ]SELECT x, y FROM t1 UNION DISTINCT SELECT x, y FROM t2; SELECT x, y FROM t1 UNION SELECT x, y FROM t2; SELECT x, y FROM t1 UNION ALL SELECT x, y FROM t2;INTERSECTINTERSECT/INTERSECT DISTINCT/INTERSECT ALL返回两侧都能找到的行交集INTERSECT与INTERSECT DISTINCT只返回去重后的行INTERSECT ALL不去重。query { INTERSECT [ ALL | DISTINCT ] } query [ .. ]SELECT x, y FROM t1 INTERSECT DISTINCT SELECT x, y FROM t2; SELECT x, y FROM t1 INTERSECT SELECT x, y FROM t2; SELECT x, y FROM t1 INTERSECT ALL SELECT x, y FROM t2;EXCEPT / MINUSEXCEPT/EXCEPT DISTINCT/EXCEPT ALL返回在左侧但不在右侧的行差集EXCEPT与EXCEPT DISTINCT只返回去重后的行EXCEPT ALL不去重MINUS是EXCEPT的同义词。query { EXCEPT [ ALL | DISTINCT ] } query [ .. ]SELECT x, y FROM t1 EXCEPT DISTINCT SELECT x, y FROM t2; SELECT x, y FROM t1 EXCEPT SELECT x, y FROM t2; SELECT x, y FROM t1 EXCEPT ALL SELECT x, y FROM t2;Lateral View 详解Lateral view子句与用户自定义表生成函数UDTF如explode()配合使用。UDTF 对每个输入行生成零行或多行输出。Lateral view 首先对基础表的每一行应用 UDTF然后将输出行与输入行连接形成具有指定表别名的虚拟表。lateralView: LATERAL VIEW [ OUTER ] udtf( expression ) tableAlias AS columnAlias [, ... ] fromClause: FROM baseTable lateralView [, ... ]列别名可以省略此时别名继承自 UDTF 返回的StructObjectInspector的字段名。参数Lateral View Outer用户可以指定可选的OUTER关键字即使通常LATERAL VIEW不会生成行例如被展开的列为空时 UDTF 不产生任何行源行将不会出现在结果中也仍生成行。使用OUTER后来自 UDTF 的列将以NULL值填充并保留源行。Multiple Lateral Views一个FROM子句可以有多个LATERAL VIEW子句。后续LATERAL VIEW可以引用其左侧出现的任何表的列。示例假设有表CREATE TABLE pageAds(pageid string, addid_list arrayint);表中包含两行数据front_page, [1, 2, 3]; contact_page, [3, 4, 5];使用LATERAL VIEW将列addid_list转换为独立行SELECT pageid, adid FROM pageAds LATERAL VIEW explode(adid_list) adTable AS adid; -- 结果 front_page, 1 front_page, 2 front_page, 3 contact_page, 3 contact_page, 4 contact_page, 5使用多个 lateral view 同时展开两列CREATE TABLE t1(c1 arrayint, c2 arrayint); SELECT myc1, myc2 FROM t1 LATERAL VIEW explode(c1) myTable1 AS myc1 LATERAL VIEW explode(c2) myTable2 AS myc2;当 UDTF 不产生行时LATERAL VIEW不会产生行可用LATERAL VIEW OUTER仍然产生行并以NULL填充对应列SELECT * FROM t1 LATERAL VIEW OUTER explode(array()) C AS a;Window Functions窗口函数详解窗口函数是对一组行称为窗口进行聚合的函数它基于行组为每一行返回聚合值。语法window_function OVER ( [ { PARTITION | DISTRIBUTE } BY colName ( [, ... ] ) ] { ORDER | SORT } BY expression [ ASC | DESC ] [ NULLS { FIRST | LAST } ] [ , ... ] [ window_frame ] )支持的窗口函数Windowing 函数LEAD、LAG、FIRST_VALUE、LAST_VALUE注意FIRST_VALUE/LAST_VALUE目前尚不支持通过参数控制跳过或保留 null 值它们总是跳过 null 值。Analytic 函数RANK、ROW_NUMBER、DENSE_RANK、CUME_DIST、PERCENT_RANK、NTILE聚合函数COUNT、SUM、MIN、MAX、AVGwindow_framewindow_frame用于指定窗口从哪一行开始、到哪一行结束支持以下格式(ROWS | RANGE) BETWEEN (UNBOUNDED | [num]) PRECEDING AND ([num] PRECEDING | CURRENT ROW | (UNBOUNDED | [num]) FOLLOWING) (ROWS | RANGE) BETWEEN CURRENT ROW AND (CURRENT ROW | (UNBOUNDED | [num]) FOLLOWING) (ROWS | RANGE) BETWEEN [num] FOLLOWING AND (UNBOUNDED | [num]) FOLLOWING默认窗口边界规则指定了ORDER BY但缺少window_frame时窗口默认为RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROWORDER BY与window_frame都缺失时窗口默认为ROW BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING。注意窗口函数中暂不支持DISTINCT。示例-- PARTITION BY 单个分区列无 ORDER BY 与窗口规格 SELECT a, COUNT(b) OVER (PARTITION BY c) FROM t; -- PARTITION BY 两个分区列无 ORDER BY 与窗口规格 SELECT a, COUNT(b) OVER (PARTITION BY c, d) FROM t; -- PARTITION BY 两个分区列带 ORDER BY SELECT a, SUM(b) OVER (PARTITION BY c, d ORDER BY e, f) FROM t; -- PARTITION BY ORDER BY 窗口规格 SELECT a, SUM(b) OVER (PARTITION BY c ORDER BY d ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) FROM t; SELECT a, AVG(b) OVER (PARTITION BY c ORDER BY d ROWS BETWEEN 3 PRECEDING AND CURRENT ROW) FROM t; SELECT a, AVG(b) OVER (PARTITION BY c ORDER BY d ROWS BETWEEN 3 PRECEDING AND 3 FOLLOWING) FROM t; SELECT a, AVG(b) OVER (PARTITION BY c ORDER BY d ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING) FROM t;Sub-Queries子查询详解FROM 子句中的子查询Hive 方言支持FROM子句中的子查询。子查询必须给定名称因为FROM子句中每个表都必须有名字子查询选择列表中的列名必须唯一这些列对外层查询如同表的列一样可用。子查询也可以是带UNION的查询表达式Hive 方言支持任意层级的子查询。select_statement FROM ( select_statement ) [ AS ] nameSELECT col FROM ( SELECT ab AS col FROM t1 ) t2WHERE 子句中的子查询Hive 方言在WHERE子句中也支持部分类型的子查询select_statement FROM table WHERE { colName { IN | NOT IN } | NOT EXISTS | EXISTS } ( subquery_select_statement )SELECT * FROM t1 WHERE t1.x IN (SELECT y FROM t2); SELECT * FROM t1 WHERE EXISTS (SELECT y FROM t2 WHERE t1.x t2.x);CTE公共表表达式详解公共表表达式CTE是在紧邻SELECT或INSERT关键字的WITH子句中指定的查询派生的临时结果集。CTE 只在单个语句的执行范围内定义并可在该范围内引用。语法withClause: WITH cteClause [ , ... ] cteClause: cte_name AS (select statement)注意WITH子句不支持在 Sub-Query 块内部使用CTE 支持在 View、CTASCREATE TABLE AS和INSERT语句中使用不支持递归查询。示例WITH q1 AS ( SELECT key FROM src WHERE key 5) SELECT * FROM q1; -- 链式 CTE WITH q1 AS ( SELECT key FROM q2 WHERE key 5), q2 AS ( SELECT key FROM src WHERE key 5) SELECT * FROM (SELECT key FROM q1) a; -- insert 场景 WITH q1 AS ( SELECT key, value FROM src WHERE key 5) FROM q1 INSERT OVERWRITE TABLE t1 SELECT *; -- CTAS 场景 CREATE TABLE t2 AS WITH q1 AS ( SELECT key FROM src WHERE key 4) SELECT * FROM q1;Transform 详解TRANSFORM子句允许用户使用指定的命令或脚本来转换输入数据。语法query: SELECT TRANSFORM ( expression [ , ... ] ) [ inRowFormat ] [ inRecordWriter ] USING command_or_script [ AS colName [ colType ] [ , ... ] ] [ outRowFormat ] [ outRecordReader ] rowFormat : ROW FORMAT (DELIMITED [FIELDS TERMINATED BY char] [COLLECTION ITEMS TERMINATED BY char] [MAP KEYS TERMINATED BY char] [ESCAPED BY char] [LINES SEPARATED BY char] | SERDE serde_name [WITH SERDEPROPERTIES property_nameproperty_value, property_nameproperty_value, ...]) outRowFormat : rowFormat inRowFormat : rowFormat outRecordReader : RECORDREADER className inRecordWriter: RECORDWRITER record_write_class注意MAP ..与REDUCE ..是 Hive 方言中对SELECT TRANSFORM ( ... )的语法等价形式因此可以用MAP/REDUCE替换SELECT TRANSFORM。参数inRowFormat指定以什么行格式将输入数据喂给运行中的脚本。默认情况下列会被转换为STRING并用TAB分隔所有NULL值会被转换为字面量字符串\N以区分NULL与空字符串。outRowFormat指定以什么行格式读取运行中脚本的输出。默认情况下用户脚本的标准输出被当作 TAB 分隔的STRING列任何只包含\N的单元格会被重新解释为NULL随后结果STRING列会按常规方式转换为表声明中指定的数据类型。inRecordWriter指定写入输入数据所用的 writer全限定类名默认为org.apache.hadoop.hive.ql.exec.TextRecordWriter。outRecordReader指定读取输出数据所用的 reader全限定类名默认为org.apache.hadoop.hive.ql.exec.TextRecordReader。command_or_script指定处理数据的命令或脚本路径。注意目前尚不支持先添加脚本文件再用脚本转换输入使用的脚本必须是本地脚本并且集群中所有主机都应能访问该脚本。colType指定命令/脚本输出应转换为的数据类型默认为STRING。AS 子句的行为对于( AS colName ( colType )? [, ... ] )?子句需要注意若实际输出列数少于用户指定的输出列数多出的用户指定输出列将填充NULL若实际输出列数多于用户指定的输出列数实际输出将被截断只保留对应列若未指定( AS colName ( colType )? [, ... ] )?子句默认输出 schema 为(key: STRING, value: STRING)key列包含第一个 TAB 之前的所有字符value列包含第一个 TAB 之后的剩余字符如果没有 TAB第二列value返回NULL。注意这与显式指定AS key, value不同——显式指定时若存在多个 TABvalue只包含第一个 TAB 与第二个 TAB 之间的部分。示例CREATE TABLE src(key string, value string); -- 基础 transform SELECT TRANSFORM(key, value) USING cat from t1; -- 指定 record writer 与 record reader SELECT TRANSFORM(key, value) ROW FORMAT SERDE MySerDe WITH SERDEPROPERTIES (p1v1,p2v2) RECORDWRITER MyRecordWriter USING cat ROW FORMAT DELIMITED FIELDS TERMINATED BY , RECORDREADER MyRecordReader FROM src; -- 使用关键字 MAP 替代 TRANSFORM FROM src INSERT OVERWRITE TABLE dest1 MAP src.key, CAST(src.key / 10 AS INT) USING cat AS (c1, c2); -- 指定 transform 输出 SELECT TRANSFORM(key, value) USING cat AS c1, c2; SELECT TRANSFORM(key, value) USING cat AS (c1 INT, c2 INT);Table Sample 详解TABLESAMPLE语句用于对表进行行采样。语法TABLESAMPLE ( num_rows ROWS )注意目前只支持采样指定数量的行。参数num_rows ROWSnum_rows是正整数常量指定采样多少行。示例SELECT * FROM src TABLESAMPLE (5 ROWS)端到端实战用 Hive 方言执行查询下面是在 Flink SQL Client 中启用 Hive 方言并执行查询的完整会话来自 Queries Overview 文档 的示例。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 2021-04-22 16:13:57,005 INFO org.apache.hadoop.mapred.FileInputFormat [] - Total input paths to process : 1 ------------ | key | value | ------------ | 1 | a | | 1 | a | | 5 | e | | 2 | b | | 3 | c | | 3 | c | | 3 | c | | 4 | d | ------------ Received a total of 8 rows从示例可以看出几个关键点UDTF 调用select explode(array(1,2,3))直接以 Hive 方言调用 Hive 的内置 UDTF这是 Lateral View 能力的基础CLUSTER BY 语义select * from tbl cluster by key的结果中同一key值的行如1、3聚在相邻分区内但整体如1与5之间并非全局有序——这正是CLUSTER BY只保证分区内有序的体现批模式要求Sort/Cluster/Distributed BY等语法主要在批模式下使用示例中显式执行了set execution.runtime-modebatch。重要提示Hive 方言不再支持 Flink SQL 查询语法。如需使用 Flink 原生语法编写查询请切换回默认方言SET table.sql-dialect default。源码视角Hive 方言查询是如何解析与执行的Hive 方言的解析能力由flink-connector-hive模块中的 parser 实现核心类位于 flink-connectors/flink-connector-hive/src/main/java/org/apache/flink/table/planner/delegation/hive/ 目录。Parser 工厂与方言注册HiveParserFactory.java 实现了ParserFactory接口其factoryIdentifier()返回SqlDialect.HIVE.name().toLowerCase()即hive并通过 SPI 机制注册到 Flink 的解析器体系中。工厂不要求任何必选/可选配置项create()方法将上下文强转为CalciteContext并构造HiveParser——注释明确说明这是因为 Hive parser 需要CalciteContext来构建 Calcite 的RelNode。HiveParser借助 Hive 自身解析器HiveParser.java 是方言解析的入口类注释为A Parser that uses Hives planner to parse a statement即直接复用 Hive 的解析器HiveASTParser生成 AST。从源码结构可以看到它维护了一组DDL_NODES常量集合覆盖TOK_CREATETABLE、TOK_DROPTABLE、TOK_ALTERTABLE、TOK_SHOWDATABASES、TOK_SWITCHDATABASE等各类 Hive DDL 与 Show 语句的 AST 节点类型用于在解析阶段区分 DDL 与 DQL 走不同的处理路径而对SELECT、INSERT等 DQL 语句则走 Hive 的查询解析与 Calcite RelNode 转换流程最终交给 Flink 执行。这解释了为何方言切换后可以零成本获得 Hive 的语法兼容性Hive 方言本质上不是重新实现一套 SQL 解析器而是把 Hive 的 parser/planner 嵌入 Flink 的执行管线。集成测试查询兼容性的验证证据仓库中提供了专门的查询兼容性集成测试 HiveDialectQueryITCase.java类注释为 Test hive query compatibility共 1149 行覆盖了大量 HiveQL 查询场景。测试在HiveCatalog上建立TableEnvironment通过HiveModuleCoreModule组合加载函数并创建foo、bar、baz、src、srcpart分区表等测试表。与之配套的还有 HiveDialectAggITCase.java聚合场景与 HiveDialectQueryPlanTest.java执行计划校验。如果你要验证本文涉及的各类查询语法在真实 Flink 环境中的行为这些测试文件是最好的参考样例。常见注意事项汇总方言互斥Hive 方言下不能使用 Flink SQL 原生查询语法如 Flink 特定的窗口函数语法需要切换回default方言。运行模式Hive 方言以批模式为主Sort/Cluster/Distribute By、Transform等语法未在流模式支持流作业请谨慎选择方言。标识符只支持db.table两级标识符不支持 Catalog 前缀。排序语义区分ORDER BY全局有序但单 task 排序可能极慢SORT BY/CLUSTER BY只保证分区内有序适合大规模数据场景。窗口函数限制FIRST_VALUE/LAST_VALUE不支持跳过/保留 null 参数控制总是跳过 null窗口函数中不支持DISTINCT。Transform 限制不支持先添加脚本文件再使用脚本必须为本地脚本且所有集群主机可访问未指定AS子句时默认输出 schema 为(key STRING, value STRING)其语义与显式AS key, value不同。CTE 限制WITH子句不支持嵌套在子查询块内不支持递归查询。Table Sample目前仅支持TABLESAMPLE (num_rows ROWS)按行数采样尚不支持按百分比、按桶等方式采样。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Flink Hive Dialect 查询语法完全指南从 SELECT 到 CLUSTER BY 的 HiveQL 实战解析Flink Hive Dialect 查询语法完全指南从 SELECT 到 CLUSTER BY 的 HiveQL 实战解析 Flink 的 Hive Dia大数据流处理批处理数据工程Flink Hive 方言子查询Sub-Queries完全指南FROM 子句与 WHERE 子句的 IN/EXISTS 用法与实现原理Flink Hive 方言子查询Sub Queries完全指南FROM 子句与 WHERE 子句的 IN/EXISTS 用法与实现原理 本指南聚焦 Apa大数据流处理批处理数据工程想提升3D打印质量Klipper固件实战指南帮你实现完美打印想提升3D打印质量Klipper固件实战指南帮你实现完美打印 你是否在为3D打印机的振纹、层错位和精度问题而烦恼Klipper固件或许就是你一直在寻找的解决嵌入式智能硬件工业制造创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考