- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】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 BY
- Group By
- Join
- Set Operation(UNION / INTERSECT / EXCEPT)
- Lateral View
- Window Functions
- Sub-Queries
- CTE(公共表表达式)
- Transform
- Table 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 Gateway(HiveServer2 Endpoint)与 Table API
- SQL 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 SETS
GROUPING 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 表示该列不属于分组集合。
ROLLUP
ROLLUP是常见分组集合类型的简写,表示给定表达式列表及该列表的所有前缀,包括空列表。例如:
GROUP BY a, b, c WITH ROLLUP等价于:
GROUP BY a, b, c GROUPING SETS ( (a, b, c), (a, b), (a), ( ) )CUBE
CUBE同样是常见分组集合类型的简写,表示给定列表及其所有可能子集——即幂集。例如:
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 BY
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;CLUSTER BY
CLUSTER 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 JOIN |
RIGHT JOIN | 返回右表全部行与左表匹配值,无匹配则追加NULL;等价于RIGHT OUTER JOIN |
FULL JOIN | 返回两侧全部行,某一侧无匹配则追加NULL;等价于FULL OUTER JOIN |
LEFT 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三种操作。
UNION
UNION/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;INTERSECT
INTERSECT/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 / MINUS
EXCEPT/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 array<int>);表中包含两行数据:
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 array<int>, c2 array<int>); 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、AVG
window_frame
window_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 ROW; ORDER 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 a+b 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、
CTAS(CREATE 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_name=property_value, property_name=property_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不同——显式指定时,若存在多个 TAB,value只包含第一个 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 ('p1'='v1','p2'='v2') 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 ROWS:
num_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-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 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-mode=batch。
重要提示: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 Hive's 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,通过HiveModule+CoreModule组合加载函数,并创建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
相关推荐
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),仅供参考