Flink—Sql接口
官网教程:https://nightlies.apache.org/flink/flink-docs-release-1.17/zh/docs/dev/table/sql/overview/
1、概念同步
1.1、逻辑层级
flink 的 catalog、database、table 和 关系型数据库的对照关系;
| 层级 | 作用 | Flink | Oracle | Greenplum | MySQL |
|---|---|---|---|---|---|
| 第一层(实例) | 最大隔离级别 | Catalog | 数据库实例(Instance) | Database | 数据库实例(Instance) |
| 第二层(命名空间) | 表的集合、schema | Database | User / Schema | Schema | Database / Schema |
| 第三层 | 数据表 | Table | Table | Table | Table |
1.2、元数据管理
flink 既然 catalog、database、table、view 的概念,那么一定有元数据需要存储;
Flink 建的 Catalog、Database、Table 元数据存在哪,取决于用的是什么 Catalog!
- 默认 catalog
default catalog:元数据存储在内存中,重启flink集群,元数据全部会丢失;
- 持久化 catalog
(Hive Catalog / Jdbc Catalog):存在外部存储(Hive Metastore / 数据库),重启不丢失,永久保存!
2、环境规划
查看 flink 中有哪些 catalog、database以及当前 catalog、database操作;
2.0、启动sql客户端
./bin/sql-client.sh
说明:该命令会启动 sql 的交互客户端界面;
2.1、catalog
- 查看
Flink SQL> show catalogs;
+-----------------+
| catalog name |
+-----------------+
| default_catalog |
+-----------------+
1 row in set
Flink SQL> show current catalog;
+----------------------+
| current catalog name |
+----------------------+
| default_catalog |
+----------------------+
1 row in set
Flink SQL>
- 切换 catalog
use catalog ${catalog_name};
说明:和切换 database 的区别是,切换 catalog 需要加上 catalog 关键字;
2.2、database
- 查看 database
Flink SQL> show databases;
+------------------+
| database name |
+------------------+
| default_database |
+------------------+
1 row in set
Flink SQL> show current database;
+-----------------------+
| current database name |
+-----------------------+
| default_database |
+-----------------------+
1 row in set
Flink SQL>
- 查询 database下的表
show tables from ${database_name};
show tables in ${database_name};
show tables from ${database_name} like '%${table_name}%';
show tables in ${database_name} like '%${table_name}%';
说明:此处查询指定的 database 可以使用 from 也可以使用 in ;
- 创建/删除 database
create database [if not exists] ${database_name};
drop database [if exists] ${database_name} [(RESTRICT | CASCADE)];
说明:删除非空数据库时,默认是 restrict,即发生异常;或者加cascade 级联删除;
- 切换 database
use ${database_name};
2.3、连接数据库
Flink 本身不自带 MySQL、Oracle 等第三方数据库的驱动包,只提供通用的 JDBC 连接器。连接具体数据库时,必须手动引入对应厂商的驱动 Jar 包,否则会报
ClassNotFoundException: com.mysql.cj.jdbc.Driver一类的错误。
连接Oracle
连接 Oracle数据库,连接器需要是 jdbc 类型,同时要把 Oracle的 驱动jar包放到 $FLINK_HOME/lib 目录下;ojdbc8.jarflink-connector-jdbc-3.1.2-1.17.jar
3、SQL 语法
- select 表
说明:这里遇到几个麻烦的细节;
1)Oracle的 date 本身自带时间,对应到 flink 需要使用 timestamp(3),不能是 timestamp(0);
2)在 flink SQL 里好像不能使用 limit N 这种记录限制,否则会报错 SQL 语句未正常结束;
- insert table
3.1、DDL相关
flnk sql 窗口查看表结构:
desc|describe ${database_name}.${table_name};
show create table ${database_name}.${table_name};
说明:以上两种都可以,desc 列出的是表的字段信息;show create table 展示的是表的完整的DDL语句;
create table
CREATE TABLE [IF NOT EXISTS] [catalog_name.][db_name.]table_name
(
{ <physical_column_definition> | <metadata_column_definition> | <computed_column_definition> }[ , ...n]
[ <watermark_definition> ]
[ <table_constraint> ][ , ...n]
)
[COMMENT table_comment]
[PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]
WITH (key1=val1, key2=val2, ...)
[ LIKE source_table [( <like_options> )] | AS select_query ]
<physical_column_definition>:
column_name column_type [ <column_constraint> ] [COMMENT column_comment]
<column_constraint>:
[CONSTRAINT constraint_name] PRIMARY KEY NOT ENFORCED
<table_constraint>:
[CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED
<metadata_column_definition>:
column_name column_type METADATA [ FROM metadata_key ] [ VIRTUAL ]
<computed_column_definition>:
column_name AS computed_column_expression [COMMENT column_comment]
<watermark_definition>:
WATERMARK FOR rowtime_column_name AS watermark_strategy_expression
<source_table>:
[catalog_name.][db_name.]table_name
<like_options>:
{
{ INCLUDING | EXCLUDING } { ALL | CONSTRAINTS | PARTITIONS }
| { INCLUDING | EXCLUDING | OVERWRITING } { GENERATED | OPTIONS | WATERMARKS }
}[, ...]
- 建表有水印
create table test_database.source_table
(
log_id bigint,
batch_id string,
job_name string,
all_flag int ,
process_id string,
etl_begin_date timestamp(3),
etl_end_date timestamp(3),
memo string,
WATERMARK FOR etl_begin_date AS etl_begin_date - INTERVAL '1' MINUTE
)
with
(
'connector' = 'jdbc',
'url' = 'jdbc:oracle:thin:@//xx.xx.xx.xx:port/sid',
'table-name' = 'schema_name.table_name', -- 必须写 模式.表名
'username' = 'xxxx',
'password' = 'xxxx',
'driver' = 'oracle.jdbc.OracleDriver'
);
说明:WATERMARK FOR etl_begin_date AS etl_begin_date - INTERVAL '1' MINUTE 水印按需指定;
- kafka消息时间
CREATE TABLE MyTable (
`user_id` BIGINT,
`name` STRING,
`record_time` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' -- reads and writes a Kafka record's timestamp
) WITH (
'connector' = 'kafka'
...
);
说明:record_time字段非业务字段,而是基于kafka消息自带的时间戳生成;kafka 自带消息时间戳有两种来源,其一为 client客户端生成数据时的时间戳(默认),另外一种为消息写入 broken时的时间戳;
- 虚拟列 virtual
CREATE TABLE MyTable (
`timestamp` BIGINT METADATA, -- part of the query-to-sink schema
`offset` BIGINT METADATA VIRTUAL, -- not part of the query-to-sink schema
`user_id` BIGINT,
`name` STRING,
) WITH (
'connector' = 'kafka'
...
);
说明:在数据从kafka 读取到 table_query 时,有 虚拟列 offset;基于 table_query 结果写入到 sink 时,不会带虚拟列 offset;
- 计算列
CREATE TABLE MyTable (
`user_id` BIGINT,
`price` DOUBLE,
`quantity` DOUBLE,
`cost` AS price * quanitity, -- evaluate expression and supply the result to queries
) WITH (
'connector' = 'kafka'
...
);
说明:同上述虚拟列,计算列在查询时可见,在 slink 端时不会进行持久化输出;
- 水印 watermark
说明:
1)水印列可以基于虚拟列或者计算列生成;
2)生成水印 watermark的列必须是不可为空的列;
- like 建表
CREATE TABLE Orders (
`user` BIGINT,
product STRING,
order_time TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TABLE Orders_with_watermark (
-- 添加 watermark 定义
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
-- 改写 startup-mode 属性
'scan.startup.mode' = 'latest-offset'
)
LIKE Orders;
| 合并策略 | 行为描述 |
|---|---|
| INCLUDING | 新表包含源表(source table)所有的表属性 ,如果与源表存在重复 key 的属性,直接失败 |
| EXCLUDING | 新表不包含源表指定的任何表属性 |
| OVERWRITING | 新表包含源表的表属性 ,但如果出现重复项,则会用新表的表属性覆盖源表中的重复表属性; |
说明:如果未提供 like 配置项(like options),默认将使用 的合并策略。INCLUDING ALL OVERWRITING OPTIONS;
alter table
- 语法格式
ALTER TABLE [IF EXISTS] table_name {
ADD { <schema_component> | (<schema_component> [, ...]) }
| MODIFY { <schema_component> | (<schema_component> [, ...]) }
| DROP {column_name | (column_name, column_name, ....) | PRIMARY KEY | CONSTRAINT constraint_name | WATERMARK}
| RENAME old_column_name TO new_column_name
| RENAME TO new_table_name
| SET (key1=val1, ...)
| RESET (key1, ...)
}
<schema_component>:
{ <column_component> | <constraint_component> | <watermark_component> }
<column_component>:
column_name <column_definition> [FIRST | AFTER column_name]
<constraint_component>:
[CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED
<watermark_component>:
WATERMARK FOR rowtime_column_name AS watermark_strategy_expression
<column_definition>:
{ <physical_column_definition> | <metadata_column_definition> | <computed_column_definition> } [COMMENT column_comment]
<physical_column_definition>:
column_type
<metadata_column_definition>:
column_type METADATA [ FROM metadata_key ] [ VIRTUAL ]
<computed_column_definition>:
AS computed_column_expression
-
rename table
alter table $old_table_name rename to $new_table_name; -
rename column
alter table ${table_name} rename ${old_column} to ${new_column}; -
add column
alter table ${table_name} add ${column_name} datatype comment '${col_comments}';
alter table ${table_name} add (
${column_a} datatype comment '${col_comments}',
${column_b} datatype comment '${col_comments}',
......
${column_n} datatype comment '${col_comments}',
);
说明:此处flink支持多字段批量添加;
-
add primary key
alter table ${table_name} add primary key(${column_name} not enforced); -
add watermark
alter table ${table_name} add watermark for ${column_name} as watermark_expression -
modify column
在这里插入代码片
- drop column
alter table ${table_name} drop $column_name};
alter table ${table_name} drop ($column_1,$column_2,...$column_n);
-
drop primary key
alter table ${table_name} drop primary key; -
drop watermark
alter table ${table_name} drop watermark; -
set / reset
-- set 'rows-per-second'
ALTER TABLE DataGenSource SET ('rows-per-second' = '10');
-- reset 'rows-per-second' to the default value
ALTER TABLE DataGenSource RESET ('rows-per-second');
说明:SET 为指定的表设置一个或多个属性。若个别属性已经存在于表中,则使用新值覆盖旧值。RESET为指定的表重置一个或多个属性。
操作SQL
ANALYSE
- 语法格式
ANALYZE TABLE [catalog_name.][db_name.]table_name PARTITION(partcol1[=val1] [, partcol2[=val2], ...]) COMPUTE STATISTICS [FOR COLUMNS col1 [, col2, ...] | FOR ALL COLUMNS]
说明:
1)对于分区表, 语法中 PARTITION(partcol1[=val1] [, partcol2[=val2], …]) 是必须指定的;
2)语法中,FOR COLUMNS col1 [, col2, …] 或者 FOR ALL COLUMNS 是可选的;
- 实际SQL
# 动态分区
ANALYZE TABLE Orders PARTITION(sold_year, sold_month, sold_day) COMPUTE STATISTICS FOR COLUMNS amount, product;
# 全表收集
ANALYZE TABLE Orders PARTITION(sold_year, sold_month, sold_day) COMPUTE STATISTICS;
explain plan
show
| 命令 | 说明 |
|---|---|
SHOW CATALOGS |
列出当前 Flink 会话中所有可用的 Catalog |
SHOW CURRENT CATALOG |
显示当前正在使用的 Catalog 名称 |
SHOW DATABASES |
列出当前 Catalog 下的所有数据库 |
SHOW CURRENT DATABASE |
显示当前正在使用的数据库名称 |
SHOW TABLES ...... |
列出当前数据库中的所有表 |
SHOW CREATE TABLE ${table_name} |
显示指定表的创建 DDL 语句 |
SHOW COLUMNS ...... |
列出指定表的所有列信息 |
SHOW VIEWS |
列出当前数据库中的所有视图 |
SHOW CREATE VIEW ${view_name} |
显示指定视图的创建 DDL 语句 |
SHOW FUNCTIONS |
列出当前会话中所有可用的函数(包括系统函数和用户自定义函数) |
SHOW MODULES |
列出当前系统中已加载的模块名称 |
SHOW FULL MODULES |
列出所有已加载模块的详细信息(包括模块名称、使用状态等) |
SHOW JARS |
列出通过 ADD JAR 命令添加的所有 JAR 文件 |
-
show table
SHOW TABLES [ ( FROM | IN ) [catalog_name.]database_name ] [ [NOT] LIKE <sql_like_pattern> ] -
show coumns
SHOW COLUMNS ( FROM | IN ) [[catalog_name.]database.]<table_name> [ [NOT] LIKE <sql_like_pattern>] -
SHOW JOBS
说明:当前 SHOW JOBS 命令只能在 SQL CLI 或者 SQL Gateway 中使用;
use
3.2、DML相关
Insert table
- select 语法
Insert { INTO | OVERWRITE } [catalog_name.][db_name.]table_name [PARTITION part_spec]
select_statement
part_spec:
(part_col_name1=val1 [, part_col_name2=val2, ...])
- values 语法
Insert { INTO | OVERWRITE } [catalog_name.][db_name.]table_name
VALUES values_row [, values_row ...]
values_row:
: (val1 [, val2, ...])
- statement set
EXECUTE STATEMENT SET
BEGIN
insert_statement;
...
insert_statement;
END;
insert_statement:
<insert_from_select>|<insert_from_values>
说明:STATEMENT SET 可以实现通过一个语句插入数据到多个表,此处的 execute 可加可不加,效果一样。
优势:
✅ 一次读取:公共数据源只扫描一次
✅ 结果复用:中间计算结果在多个输出间共享
✅ 单作业提交:避免多作业的调度开销
3.3、DQL相关
with 语法
with 语法类似于 关系型数据库的 with 视图,语法格式完全一样;
select 语法
- 常规语法
select col_1,col_2 from ${table_name}/${view_name};
- 特殊语法
SELECT order_id, price FROM (VALUES (1, 2.0), (2, 3.1)) AS t (order_id, price);
说明:查询操作还可以在 VALUES 子句中使用内联数据。每一个元组对应一行,另外可以通过设置别名来为每一列指定名称。
select 'hello world!';
select current_timestamp;
# 查看可用函数
show functions;
- 时态关联
SELECT o.order_id, o.total, c.country, c.zip
FROM Orders AS o -- 订单流(事实表)
inner JOIN Customers
FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.customer_id = c.id; -- 关联条件
window 语法
滚动窗口
- 语法格式
TUMBLE(TABLE data, DESCRIPTOR(timecol), size [, offset ])
参数说明:
1)TUMBLE函数需要三个必要的参数,其中一个可选参数:
2)该函数主要用来将业务时间按照固定的窗口大小进行统计;
3)offset 参数的作用是调整窗口的其实值,在默认值的基础上进行偏移;
例如,一个10分钟窗口(INTERVAL ‘10’ MINUTE),默认的窗口是 00:00-00:10, 00:10-00:20… 如果加一个 offset 为 INTERVAL ‘4’ MINUTE,窗口就会变成 00:04-00:14, 00:14-00:24… 。
- 无 offset 语法
select job_name,
batch_id,
etl_begin_date,
window_start,
window_end,
window_time
from table(tumble(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '1' hours
)
);
- 有 offset 语法
select job_name,
batch_id,
etl_begin_date,
window_start,
window_end,
window_time
from table(tumble(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '1' hours,
interval '10' minutes
)
);
- group by
select window_start, window_end, window_time, count(0) as jls
from table(tumble(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '1' hours
)
) h
group by window_start, window_end, window_time;
滑动窗口
-
语法格式
HOP(TABLE data, DESCRIPTOR(timecol), slide, size [, offset ])
说明:HOP需要四个必要参数,其中一个可选参数; -
明细查询
select job_name,
batch_id,
etl_begin_date,
window_start,
window_end,
window_time
from table(hop(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '30' minutes,
interval '1' hours
)
);
说明:窗口大小为 1 小时,滑动频率为 20 分钟;
- group by
select window_start, window_end, window_time, count(0) as jls
from table(hop(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '30' minutes,
interval '1' hours
)
) h
group by window_start, window_end, window_time;
累计窗口
-
语法格式
CUMULATE(TABLE data, DESCRIPTOR(timecol), step, size)
说明:CUMULATE需要四个必要参数,其中一个可选参数; -
明细查询
select job_name,
batch_id,
etl_begin_date,
window_start,
window_end,
window_time
from table(cumulate(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '30' minutes,
interval '1' hours
)
);
- group by
select window_start, window_end, window_time, count(0) as jls
from table(cumulate(table ss_test_log_stat,
descriptor(etl_begin_date),
interval '30' minutes,
interval '1' hours
)
) h
group by window_start, window_end, window_time;
窗口聚合
没搞懂
分组聚合
没搞懂
over 聚合
- 语法格式
selct
col_m, col_n,
agg_func(agg_col) OVER (
[PARTITION BY col_1,col_2[,...] ORDER BY time_col
range_definition)
from ${tabel_name} t
说明:可以在一个 SQL 里定义多个窗口,但是每个窗口必须一致;
- 单 over 聚集
select job_name,
batch_id,
etl_begin_date,
max(batch_id) over(partition by job_name order by etl_begin_date
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW) as max_bat
from ss_test_log_stat;
- 多 over聚集
select job_name,
batch_id,
etl_begin_date,
max(batch_id) over(partition by job_name order by etl_begin_date
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW) as max_bat,
min(batch_id) over(partition by job_name order by etl_begin_date
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW) as min_bat
from ss_test_log_stat;
- 多 over 新语法
select job_name,
batch_id,
etl_begin_date,
max(batch_id) over w as max_bat,
min(batch_id) over w as min_bat
from ss_test_log_stat t
window w as (partition by job_name order by etl_begin_date
RANGE BETWEEN INTERVAL '1' HOUR PRECEDING AND CURRENT ROW);
表关联
SELECT order_id, res
FROM Orders,
LATERAL TABLE(table_func(order_id)) t(res)
窗口关联
集合操作
- union / union all
# m表和 n表结构集合并,去重
select col_name,col_code,col_n from ${table_m} m
union
select col_name,col_code,col_n from ${table_n} n
# m表和 n表结构集合并,不去重
select col_name,col_code,col_n from ${table_m} m
union all
select col_name,col_code,col_n from ${table_n} n
- intersect / intersect all
# 表m 和 表n 交集,去重
select col_name,col_code,col_n from ${table_m} m
intersect
select col_name,col_code,col_n from ${table_n} n
# 表m 和 表n 交集,不去重
select col_name,col_code,col_n from ${table_m} m
intersect all
select col_name,col_code,col_n from ${table_n} n
- except / except all
# 在表m 但是不再表n 的记录,去重
select col_name,col_code,col_n from ${table_m} m
except
select col_name,col_code,col_n from ${table_n} n
# 在表m 但是不再表n 的记录,不去重
select col_name,col_code,col_n from ${table_m} m
except
select col_name,col_code,col_n from ${table_n} n
- in / exists
select col_name, col_code, col_n
from ${table_m} m
where m.col_name in
(select col_name from ${table_n} n);
select col_name, col_code, col_n
from ${table_m} m
where m.col_name exists
(select col_name from ${table_n} n);
说明:此处 in / exists 是同样的功能,最终 flink 会将其重写为 表关联和分组方式;
Top-N查询
- 语法格式
SELECT [column_list]
FROM (
SELECT [column_list],
ROW_NUMBER() OVER ([PARTITION BY col1[, col2...]]
ORDER BY col1 [asc|desc][, col2 [asc|desc]...]) AS rownum
FROM table_name)
WHERE rownum <= N [AND conditions]
窗口去重
- 语法格式
SELECT [column_list]
FROM (
SELECT [column_list],
ROW_NUMBER() OVER (PARTITION BY window_start, window_end [, col_key1...]
ORDER BY time_attr [asc|desc]) AS rn
FROM table_name) -- relation applied windowing TVF
WHERE (rn = 1 | rn <=1 | rn < 2) [AND conditions]
说明:理论上,窗口去重是窗口顶N的一种特例,其中N为1,按处理时间或事件时间的顺序。
- 示例SQL
SELECT *
FROM (SELECT bidtime,
price,
item,
supplier_id,
window_start,
window_end,
ROW_NUMBER() OVER(PARTITION BY window_start, window_end ORDER BY bidtime DESC) AS rn
FROM TABLE(TUMBLE(TABLE Bid,
DESCRIPTOR(bidtime),
INTERVAL '10' MINUTES))) h
WHERE rn <= 1;
3.4、Sql Hint
-- 覆盖查询语句中源表的选项
select id, name from kafka_table1 /*+ OPTIONS('scan.startup.mode'='earliest-offset') */;
sink.partitioner='round-robin'
表示轮询分发,让数据均匀写入 Kafka 所有分区,解决默认分区策略可能导致的数据倾斜问题。
- BROADCAST
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Flink 会使用 broadcast join,且表 t1 会被当作需 broadcast 的表。
SELECT /*+ BROADCAST(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- Flink 会在两个联接中都使用 broadcast join,且 t1 和 t3 会被作为需 broadcast 到下游的表。
SELECT /*+ BROADCAST(t1, t3) */ * FROM t1 JOIN t2 ON t1.id = t2.id JOIN t3 ON t1.id = t3.id;
-- BROADCAST 只支持等值的联接条件
-- 联接提示会失效,只能使用支持非等值条件联接的 nested loop join。
SELECT /*+ BROADCAST(t1) */ * FROM t1 join t2 ON t1.id > t2.id;
-- BROADCAST 不支持 `Full Outer Join`
-- 联接提示会失效,planner 会根据 cost 选择最合适的联接策略。
SELECT /*+ BROADCAST(t1) */ * FROM t1 FULL OUTER JOIN t2 ON t1.id = t2.id;
注意: BROADCAST 只支持等值的联接条件,且不支持 Full Outer Join。
-
SHUFFLE_HASH
-
SHUFFLE_MERGE
-
NEST_LOOP
注意:NEST_LOOP 同时支持等值的和非等值的联接条件。
================================== over =============================================
所有评论(0)