官网教程: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.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 =============================================

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。