Flink Table 与 SQL
概述
Flink Table & SQL 是流批统一的高级 API:把数据流抽象为动态表,用 SQL 持续查询。它支持 Kafka、MySQL、CDC 等连接器,内置 Temporal Join 与维表关联。本文讲透动态表模型、持续查询、时间属性与常见 Join。
一、动态表模型
1.1 核心思想
数据流 = 不断追加的动态表
SQL 查询 = 持续查询(每有新数据就更新结果)
结果 = 动态结果表(输出为流)| 概念 | 说明 |
|---|---|
| 动态表 | 随数据变化的表(追加/更新) |
| 持续查询 | 对动态表的连续 SQL 查询 |
| 结果表 | 查询结果动态更新 |
| 转流 | 结果表按追加/更新模式输出 |
1.2 示例
sql
-- 从 Kafka 读订单流
CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10, 2),
ts TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);
-- 持续查询
SELECT user_id, SUM(amount) AS total
FROM orders
GROUP BY user_id;二、时间属性
2.1 两种时间
| 属性 | 说明 |
|---|---|
| 事件时间 | 业务时间 + 水位线 |
| 处理时间 | 系统处理时刻 |
2.2 声明
sql
-- 事件时间(定义水位线)
ts TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
-- 处理时间
proc_time AS PROCTIME()| 用途 | 说明 |
|---|---|
| 窗口聚合 | 需要时间属性 |
| 维表 Join | 处理时间临时关联 |
| 去重/排序 | 事件时间语义 |
三、查询连接器
3.1 Source/Sink 连接器
| 连接器 | 用途 |
|---|---|
| Kafka | 主流流输入输出 |
| JDBC | MySQL/PG 读写 |
| Filesystem | 文件/Parquet |
| Elasticsearch | 结果写入 |
| HBase | 维表查询 |
| CDC | 数据库变更流 |
3.2 表定义示例(JDBC Sink)
sql
CREATE TABLE result (
user_id BIGINT,
total DECIMAL(10, 2)
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://db:3306/report',
'table-name' = 'user_total',
'username' = 'report',
'password' = '***'
);
INSERT INTO result SELECT user_id, SUM(amount)
FROM orders GROUP BY user_id;四、CDC Connector
4.1 概念
CDC(Change Data Capture)连接器读取数据库变更日志,把增删改转成流:
| 连接器 | 数据源 |
|---|---|
| MySQL CDC | binlog |
| PostgreSQL CDC | WAL |
| Oracle CDC | redo log |
4.2 示例
sql
CREATE TABLE cdc_orders (
id BIGINT, user_id BIGINT, amount DECIMAL(10, 2)
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql',
'database-name' = 'orderdb',
'table-name' = 'orders',
'username' = 'flink',
'password' = '***'
);| 能力 | 说明 |
|---|---|
| 全量 + 增量 | 先快照后 binlog |
| 精确一次 | 位点记录 |
| Schema 变更 | 支持部分演化 |
五、Temporal Join
5.1 时间旅行关联
关联某个时间点的版本数据(如版本化表/变更流):
sql
SELECT o.order_id, p.price
FROM orders o
LEFT JOIN products FOR SYSTEM_TIME AS OF o.ts AS p
ON o.product_id = p.id| 要求 | 说明 |
|---|---|
| 版本表 | 带时间属性的表 |
| 等值条件 | 主键关联 |
| 语义 | 取关联时刻的版本 |
5.2 应用
| 场景 | 说明 |
|---|---|
| 价格历史 | 订单关联下单时价格 |
| 配置版本 | 关联当时生效配置 |
| 维度快照 | 快照关联 |
六、维表 Join(Lookup Join)
6.1 概念
关联外部存储(MySQL/HBase/Redis)的实时查询维表:
sql
CREATE TABLE dim_user (
user_id BIGINT PRIMARY KEY NOT ENFORCED,
city STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://db:3306/dim',
'table-name' = 'dim_user',
'lookup.cache.max-rows' = '10000',
'lookup.cache.ttl' = '60s'
);
SELECT o.order_id, u.city
FROM orders o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u
ON o.user_id = u.user_id;6.2 缓存优化
| 参数 | 说明 |
|---|---|
lookup.cache.max-rows | 缓存行数 |
lookup.cache.ttl | 缓存有效期 |
lookup.max-retries | 查询重试 |
| 策略 | 场景 |
|---|---|
| 无缓存 | 数据一致、低频 |
| 缓存 | 高频查询、容忍短时陈旧 |
6.3 与 Temporal Join 区别
| 对比 | Temporal Join | Lookup Join |
|---|---|---|
| 数据来源 | Flink 表/流 | 外部存储实时查 |
| 时间语义 | 事件时间版本 | 处理时间查询 |
| 典型 | 版本化数据 | 维表补全 |
七、持续查询与结果输出
7.1 查询类型
| 类型 | 结果模式 | 示例 |
|---|---|---|
| 追加查询 | 只追加 | 过滤、投影 |
| 更新查询 | 更新/删除 | 聚合、去重 |
| 全量更新 | 全量重算 | 无 key 聚合 |
7.2 输出模式(表转流)
| 模式 | 说明 |
|---|---|
| Append | 只发新增 |
| Retract | 先撤销再新增 |
| Upsert | 按主键更新 |
八、流批统一
| 能力 | 说明 |
|---|---|
| 同一 SQL | 流批一套语法 |
| 执行模式 | STREAMING / BATCH |
| 连接器 | 批读文件、流读 Kafka |
sql
SET 'execution.runtime-mode' = 'BATCH';
-- 或 'STREAMING'常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 水位线未推进 | 表未定义 WATERMARK |
| 维表 Join 慢 | 开缓存、减少远程查询 |
| CDC 全量阶段慢 | 大表加并行与分片 |
| 结果更新异常 | 检查主键与 Upsert 模式 |
| Temporal 报错 | 版本表需带时间属性 |