官网教程:https://nightlies.apache.org/flink/flink-docs-release-1.17/zh/docs/dev/table/sql/overview/

1、概念同步

1.1、逻辑层级

flink 的 catalog、database、table 和 关系型数据库的对照关系;

层级作用FlinkOracleGreenplumMySQL
第一层(实例)最大隔离级别Catalog数据库实例(Instance)Database数据库实例(Instance)
第二层(命名空间)表的集合、schemaDatabaseUser / SchemaSchemaDatabase / Schema
第三层数据表TableTableTableTable

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.jar
flink-connector-jdbc-3.1.2-1.17.jar

JDBC连接器下载地址:https://repo1.maven.org/maven2/org/apache/flink/flink-connector-jdbc/3.1.2-1.17/flink-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 =============================================

更多推荐