Flink Hive 方言查询(Queries)完全指南:从 SELECT 语法到 Sort/Cluster/Join/CTE 实战
2026/9/24 8:41:35 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】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 ByJoin、集合操作、Lateral View、窗口函数、子查询、CTE、TransformTable 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 支持defaulthive两种 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 BYTransform等语法尚未在流(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 子句

ALLDISTINCT选项指定是否返回重复行:

  • 两者都未指定时,默认为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 BYDISTRIBUTE BYSORT 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);

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 方言支持UNIONINTERSECTEXCEPT/MINUS三种操作。

UNION

UNION/UNION DISTINCT/UNION ALL返回两侧都能找到的行:

  • UNIONUNION 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返回两侧都能找到的行(交集):

  • INTERSECTINTERSECT 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返回在左侧但不在右侧的行(差集):

  • EXCEPTEXCEPT DISTINCT只返回去重后的行;
  • EXCEPT ALL不去重;
  • MINUSEXCEPT的同义词。
<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 函数LEADLAGFIRST_VALUELAST_VALUE
    • 注意:FIRST_VALUE/LAST_VALUE目前尚不支持通过参数控制跳过或保留 null 值,它们总是跳过 null 值。
  • Analytic 函数RANKROW_NUMBERDENSE_RANKCUME_DISTPERCENT_RANKNTILE
  • 聚合函数COUNTSUMMINMAXAVG

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 BYwindow_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 ] name
SELECT col FROM ( SELECT a+b AS col FROM t1 ) t2

WHERE 子句中的子查询

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)是在紧邻SELECTINSERT关键字的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 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-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

从示例可以看出几个关键点:

  1. UDTF 调用select explode(array(1,2,3))直接以 Hive 方言调用 Hive 的内置 UDTF,这是 Lateral View 能力的基础;
  2. CLUSTER BY 语义select * from tbl cluster by key的结果中,同一key值的行(如13)聚在相邻分区内,但整体(如15之间)并非全局有序——这正是CLUSTER BY只保证分区内有序的体现;
  3. 批模式要求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_CREATETABLETOK_DROPTABLETOK_ALTERTABLETOK_SHOWDATABASESTOK_SWITCHDATABASE等各类 Hive DDL 与 Show 语句的 AST 节点类型,用于在解析阶段区分 DDL 与 DQL 走不同的处理路径;而对SELECTINSERT等 DQL 语句,则走 Hive 的查询解析与 Calcite RelNode 转换流程,最终交给 Flink 执行。

这解释了为何方言切换后可以"零成本"获得 Hive 的语法兼容性:Hive 方言本质上不是重新实现一套 SQL 解析器,而是把 Hive 的 parser/planner 嵌入 Flink 的执行管线。

集成测试:查询兼容性的验证证据

仓库中提供了专门的查询兼容性集成测试 HiveDialectQueryITCase.java(类注释为 "Test hive query compatibility"),共 1149 行,覆盖了大量 HiveQL 查询场景。测试在HiveCatalog上建立TableEnvironment,通过HiveModule+CoreModule组合加载函数,并创建foobarbazsrcsrcpart(分区表)等测试表。与之配套的还有 HiveDialectAggITCase.java(聚合场景)与 HiveDialectQueryPlanTest.java(执行计划校验)。如果你要验证本文涉及的各类查询语法在真实 Flink 环境中的行为,这些测试文件是最好的参考样例。

常见注意事项汇总

  1. 方言互斥:Hive 方言下不能使用 Flink SQL 原生查询语法(如 Flink 特定的窗口函数语法),需要切换回default方言。
  2. 运行模式:Hive 方言以批模式为主,Sort/Cluster/Distribute ByTransform等语法未在流模式支持,流作业请谨慎选择方言。
  3. 标识符:只支持db.table两级标识符,不支持 Catalog 前缀。
  4. 排序语义区分ORDER BY全局有序但单 task 排序可能极慢;SORT BY/CLUSTER BY只保证分区内有序,适合大规模数据场景。
  5. 窗口函数限制FIRST_VALUE/LAST_VALUE不支持跳过/保留 null 参数控制(总是跳过 null);窗口函数中不支持DISTINCT
  6. Transform 限制:不支持先添加脚本文件再使用;脚本必须为本地脚本且所有集群主机可访问;未指定AS子句时默认输出 schema 为(key STRING, value STRING),其语义与显式AS key, value不同。
  7. CTE 限制WITH子句不支持嵌套在子查询块内,不支持递归查询。
  8. Table Sample:目前仅支持TABLESAMPLE (num_rows ROWS)按行数采样,尚不支持按百分比、按桶等方式采样。
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

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

相关推荐

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

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

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

立即咨询