Flink(五):Flink SQL、CDC 与生产实践
Flink(五):Flink SQL、CDC 与生产实践
导语:真正落地的 Flink 作业,90% 是 SQL。这一篇讲清 SQL 层最反直觉的三件事——动态表与"持续查询"、SQL 里的时间属性与双流 Join、维表关联的三种姿势;再展开 CDC 同步链路、数据湖集成与实时数仓分层架构,最后给出常见 SQL 性能问题与排查套路,共 13 题。
一、Flink SQL 核心模型
1. Flink SQL 的「动态表 + 持续查询」是什么意思?
答: 这是理解 Flink SQL 的唯一入口。核心一句话:
流 → 动态表 → SQL → 动态表 → 流
三个必懂的概念:
| 概念 | 含义 | 关键点 |
|---|---|---|
| 动态表(Dynamic Table) | 与静态表相对,随时间不断变化的表:流里的每条事件 = 对表的一次操作(Insert / Update / Delete) | 动态表永不"查询结束",因为流没有终点 |
| 持续查询(Continuous Query) | 一个 SQL 定义在动态表上,永不停止地随着输入变化而更新结果表 | 结果表是不断被修正的(这是"更新流"的来源) |
| 变更日志流(Changelog Stream) | 动态表与流之间的双向映射:INSERT(+I)、UPDATE_BEFORE(-U)、UPDATE_AFTER(+U)、DELETE(-D) | +I/-U/+U/-D 是 Flink SQL 输出的"四种消息类型" |
由此推出的核心结论(面试高频):
- 同一个 SQL,流模式(STREAMING)与批模式(BATCH)语义不同:流模式是「持续查询、结果不断修正」,批模式是「一次性计算、只输出最终结果」;
- 「结果表」可以是非追加的(retract/upsert):例如
count(*)的结果会随输入变化而更新,所以下游要能处理-U/+U; - 「被更新的表」不能简单写进普通 Sink:
GROUP BY/JOIN/ 去重等产生回撤消息的查询,输出到 Kafka 时需要upsert-kafka连接器或指定changelog-format; - 不是所有查询都能"更新":只有有界算子/有状态算子才能处理回撤,纯 append-only 场景才可以用普通 Sink。
试金石问题:「
SELECT count(*) FROM t GROUP BY a的输出流长什么样?」 → 每来一条数据就输出一行(含-U旧值 ++U新值),而不是"到最后才给一个总数"。
2. SQL 里怎么写时间属性(Event Time / Processing Time)和 Watermark?
答: 三种写法:
① 用 WATERMARK FOR ... AS ...(推荐,来自 DDL)
CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10,2),
order_time TIMESTAMP(3),
-- 声明为事件时间:延迟 5 秒(rows:水位 = 最大事件时间 - 5s)
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);② 用计算列(Computed Column)
-- 事件时间:直接从字段取
event_time AS TO_TIMESTAMP(FROM_UNIXTIME(ts / 1000)),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
-- 处理时间:用 PROCTIME()(会生成一个隐藏的"处理时间属性列")
proc_time AS PROCTIME()③ 用 DataStream 的 assignTimestampsAndWatermarks 后转表(不推荐,SQL 层能用就用 SQL)。
三个高频考点:
| 考点 | 说明 |
|---|---|
WATERMARK 必须定义在 TIMESTAMP(3) 或 TIMESTAMP_LTZ(3) 列上 | 且该列必须在表里真实存在(计算列也算) |
PROCTIME() 是"处理时间属性" | 它不会真正落成一个物理列,只是让窗口/时态表 join 能拿到处理时间;PROCTIME() 不能用于 Event Time 窗口 |
SQL 的窗口:TUMBLE / HOP / CUMULATE / SESSION(TVF) | TUMBLE(TABLE t, DESCRIPTOR(ts), INTERVAL '1' MINUTE) 是表值函数(TVF)写法,是 Flink 1.13+ 的推荐形态,替代了旧的 GROUP BY TUMBLE(ts, ...) |
-- 窗口 TVF(推荐写法)
SELECT window_start, window_end, user_id, SUM(amount) AS total
FROM TABLE(
TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '1' MINUTE)
)
GROUP BY window_start, window_end, user_id;特别说明 CUMULATE(累积窗口):它是 Flink 特有的窗口,用于解决「滑动窗口成本高,但业务又想看"累计到今天当前的进度"」的场景——同一批数据只在 window 边界被计算一次,比同效果的滑动窗口省得多的状态。
3. 什么是维表关联?Lookup Join、Temporal Join、Interval Join 怎么选?
答: 「流 × 维表」是实时数仓最核心的算子,Flink SQL 提供三种:
| 方式 | SQL 形态 | 语义 | 状态/存储 | 适用 |
|---|---|---|---|---|
| Lookup Join(维表点查) | LEFT JOIN dim FOR SYSTEM_TIME AS OF p.proc_time AS d ON ... | 每条数据实时查外部维表(HBase/MySQL/Redis) | 不占 Flink 状态,占外部查询压力 | 维表大、变化频繁、只需"当前值" |
| Temporal Join(时态表关联) | LEFT JOIN dim FOR SYSTEM_TIME AS OF o.order_time AS d ON ... | 按事件时间取"当时"的维表版本 | 占 Flink 状态(保留维表历史版本) | 需要"历史快照"语义(如汇率、费率) |
| Interval Join(区间 Join) | ... WHERE b.ts BETWEEN a.ts - INTERVAL '5' MINUTE AND a.ts | 两条流按时间区间关联 | 状态可回收 | 双流关联(订单 × 支付) |
Lookup Join 的关键细节(最常用,问得最多):
-- 维表 DDL:需要声明 'lookup' 类型的 connector
CREATE TABLE dim_product (
product_id BIGINT,
name STRING,
price DECIMAL(10,2),
PRIMARY KEY (product_id) NOT ENFORCED -- Lookup Join 需要主键
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://mysql:3306/dw',
'table-name' = 'dim_product',
'lookup.cache' = 'PARTIAL', -- 缓存策略
'lookup.partial-cache.max-rows' = '100000',
'lookup.partial-cache.expire-after-write' = '10min'
);
-- 查询:用 PROCTIME 触发点查
SELECT o.order_id, p.name, p.price
FROM orders AS o
LEFT JOIN dim_product FOR SYSTEM_TIME AS OF o.proc_time AS p
ON o.product_id = p.product_id;Lookup Join 的四个核心问题:
- 为什么要
FOR SYSTEM_TIME AS OF o.proc_time? → 这是在告诉 Flink「在处理的这一刻去查维表」,所以它是处理时间语义的点查(不是事件时间,事件时间要用 Temporal Join); - 缓存策略怎么选?
lookup.cache | 行为 | 取舍 |
|---|---|---|
NONE(默认) | 每条数据都查一次 | 最准,但 QPS 打到外部库,容易打爆 |
PARTIAL(推荐) | 缓存一部分(LRU + 过期),未命中再查 | 兼顾性能与新鲜度,最常用 |
FULL | 全量加载维表(要求维表有界) | 性能最好,但维表更新不及时、内存占用大 |
- 维表关联不上怎么办(
LEFT JOIN得到 null)? → ①允许延迟:用table.exec.async-lookup.output-mode/ 重试;②把 null 结果走侧输出补偿;③ 用DimJoinCache的「未命中不缓存 null」策略(否则维表补数后依然查不到); - 异步查询:
'lookup.async' = 'true'让 Lookup Join 走 Async I/O(底层就是《Flink(四)》的 Async I/O),必须开,否则每条数据同步等 DB 响应,吞吐极低。
Temporal Join 的两种实现:
- 基于 changelog 的时态表(维表本身是一条流,带
PRIMARY KEY+ Watermark)→ Flink 在状态里保留所有历史版本(所以要求维表有界或带 TTL,否则状态爆炸); - Lookup 实现的时态表(用
LOOKUPhint)→ 按事件时间去外部查「历史版本」(需要外部表能按时间回溯)。
4. SQL 的常用优化手段有哪些?
答: 按「执行计划 → 算子 → 状态」分三层:
① 计划层(最有效、最先做)
| 手段 | 说明 |
|---|---|
MiniBatch 聚合(table.exec.mini-batch.enabled) | 把「每条数据一次聚合」变成「攒一小批再聚合」,大幅减少状态访问与输出回撤(默认就开:allow-latency=5s, size=1000) |
两阶段聚合(table.optimizer.agg-phase-strategy: TWO_PHASE) | 解决聚合倾斜(类似《Flink(四)》的加盐,SQL 层开参数即可) |
| Local-Global 去重(两阶段去重) | 用 ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ...) 做去重时,开 table.optimizer.distinct-agg.split.enabled 或 table.optimizer.agg-phase-strategy 减少热点 |
| Join 顺序与 Broadcast | 小表用 BROADCAST hint(/*+ BROADCAST(dim) */)避免 Shuffle;内表 join 时 Flink 也会自动选择 |
| TopN 优化 | ROW_NUMBER() TopN 配 MiniBatch + LocalGlobal,减少状态与回撤 |
| MiniBatch Join | 减少 join 的状态访问次数 |
② 算子层
| 手段 | 说明 |
|---|---|
table.exec.source.idle-timeout | 空分区(idle source)忽略,避免 Watermark 卡死(等价于 DataStream 的 withIdleness,必配) |
| 设置算子并行度 | SET 'parallelism.default' 或 table.exec.resource.default-parallelism;也可在 SQL 里用 /*+ PARALLELISM(4) */ |
table.optimizer.non-deterministic-update.strategy | 处理非确定性更新(如用了 NOW()/随机)导致的回撤乱序 |
table.exec.sink.not-null-enforcer | 防止 null 写进 NOT NULL 列 |
③ 状态层
| 手段 | 说明 |
|---|---|
| 设置状态 TTL | SET 'table.exec.state.ttl' = '24h' —— SQL 作业最常见的内存泄漏就是忘了设它(默认是 无过期!) |
| 开启 MiniBatch + 减少回撤 | 回撤消息多 → 状态与下游压力都大 |
| RocksDB + 增量 Checkpoint | 与 DataStream 一致(SQL 作业也是同一个 StateBackend) |
一条"起手配置"(生产常用):
# flink-conf.yaml 或 SQL 会话配置
table.exec.state.ttl: 24h
table.exec.mini-batch.enabled: true
table.exec.mini-batch.allow-latency: 5s
table.exec.mini-batch.size: 2000
table.optimizer.agg-phase-strategy: TWO_PHASE
table.optimizer.distinct-agg.split.enabled: true
table.exec.source.idle-timeout: 5min重点提醒:
table.exec.state.ttl默认是"永不过期",SQL 作业上线前必须显式设置,否则一段时间后状态就会涨到 Checkpoint 超时。
二、CDC 与数据同步
5. Flink CDC 是什么?原理是什么?
答: Flink CDC 是基于「变更数据捕获(Change Data Capture)」把数据库的增量变更同步进 Flink 的连接器,实现「全量 + 增量一体化」的实时同步。
原理:
关键技术点:
| 技术点 | 说明 |
|---|---|
| 全量读的一致性 | CDC 用 FLUSH TABLES WITH READ LOCK / 事务快照 保证「全量阶段」与「增量阶段」不重不漏(读快照时的位点必须精确记录) |
| 增量读 | 通过 binlog(MySQL ROW 格式)/ WAL / redo log 解析变更 |
| Exactly-Once | 位点写入 Flink 状态并随 Checkpoint 持久化;下游用事务/幂等 |
| 断点续传 | 从 Checkpoint 恢复时直接从上次位点继续,不用重新全量 |
| 无锁/并行快照(新版) | 支持多并行度并行快照(按 chunk 切分表数据),加速大表首同步 |
版本辨析(容易答错):
- Flink CDC 2.x:基于 DataStream API(
MySqlSource、MySqlParallelSource),需要写 Java 代码; - Flink CDC 3.x:推出 CDC Pipeline(YAML 配置驱动),无需写 Java,且内置 schema evolution(表结构变更同步)——这是目前的主流形态;
- 两种接入方式:① DataStream 自定义代码(灵活);② Flink SQL 的 CDC Connector(
connector = 'mysql-cdc',最常用)。
-- SQL 方式(最简单)
CREATE TABLE mysql_orders (
id BIGINT PRIMARY KEY NOT ENFORCED,
user_id BIGINT,
amount DECIMAL(10,2),
update_time TIMESTAMP(3)
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql',
'port' = '3306',
'username' = 'flink',
'password' = '***',
'database-name' = 'order_db',
'table-name' = 'orders'
);6. 基于 Flink 的实时同步链路怎么做?有什么坑?
答: 标准链路:
七个实战坑(面试很吃这套):
| 坑 | 说明 | 对策 |
|---|---|---|
| 1. 首同步打爆源库 | 全量阶段多个并行度全表扫描 → 主库 IO/CPU 飙升 | 在从库读、限制并行度、拉长时间窗、错峰 |
| 2. binlog 保留时间不够 | 作业停太久,位点对应的 binlog 已被清理 → 无法续传 | 调大 binlog 保留时间,并监控 Lag |
| 3. 表结构变更(DDL) | 加字段/改类型后,下游 Schema 对不上 | CDC 3.x 的 schema evolution;或下游用宽松格式(如 Avro + Registry) |
| 4. 大事务 / 长事务 | 一个大事务产生海量 binlog,一次性灌入导致反压 | 拆事务、Kafka 缓冲、提升下游写能力 |
| 5. 下游写入压力 | 直接写 DB 单条写必然崩 | 批量写 + upsert + 分库分表/分区;Doris 用 Stream Load |
| 6. 数据重复 | At-Least-Once 语义下重放导致重复 | 下游 主键 upsert(upsert-kafka、Doris Unique 表) |
| 7. 时区与时间字段 | MySQL TIMESTAMP 与时区/秒精度处理不一致 | 统一用 TIMESTAMP_LTZ 或明确 server-time-zone |
7. 什么是「数据湖 + Flink」的实时湖仓架构?
答: Flink + 数据湖(Paimon / Hudi / Iceberg) 是当前实时数仓的主流方向,解决「流式写入 + 批式查询」的统一。
| 湖格式 | 特点 | 与 Flink 的关系 |
|---|---|---|
| Apache Paimon(原 Flink Table Store) | 为流式更新而生:LSM 结构、支持主键 upsert、changelog 生产、小文件合并 | 与 Flink 同源,流读流写最顺,实时湖仓首选 |
| Apache Hudi | Upsert 能力强、支持 Copy-on-Write / Merge-on-Read、增量查询 | Flink/Spark 都支持,传统离线湖仓升级路线多 |
| Apache Iceberg | 表格式标准、生态最广(Trino/Spark/Flink)、隐藏分区、Time Travel | Flink 支持写入与流读,但 upsert 能力弱于 Paimon/Hudi |
为什么需要湖仓(价值):
① 一份存储两用:实时写入、离线/交互式读取,不用维护"Kafka 明细 + 离线表"两份;
② 支持 upsert 与 Time Travel:主键更新、回溯历史、快照隔离;
③ 成本更低:对象存储(S3/OSS/HDFS)比 Kafka 长期存储便宜得多。
三、架构与选型
8. 一个完整的实时数仓分层架构怎么设计?
答: 经典 ODS → DWD → DWS → ADS 四层(与离线数仓同构,只是"实时化"):
分层设计的三个"为什么":
| 问题 | 答 |
|---|---|
| 为什么要分层,直接一路到底不行吗? | 分层让复用与隔离成为可能:DWD 一份明细可以支撑多个 DWS 主题;出问题时能定位到层;也方便"实时 + 离线共用同一份 DWD" |
| 为什么实时 DWD 要落湖? | ①回刷能力(口径变化可重算);②明细留存(实时表通常只留聚合结果);③离线/实时口径统一(同一份数据两种计算引擎读) |
| 为什么 ADS 用 OLAP 而不是直接查 Paimon? | ADS 面向高并发点查与看板,Doris/ClickHouse/Redis 在并发与查询延迟上更合适;湖表更适合大范围扫描与回溯 |
9. Flink 实时链路与离线链路怎么保持一致(口径统一)?
答: 这是实时数仓最大的工程难题,面试问「一致性」时要从这四个层面答:
| 层面 | 问题 | 做法 |
|---|---|---|
| 数据源一致 | 实时读 Kafka、离线读 Hive,上游是否同一份? | 统一用 CDC 落湖,实时/离线都从同一份 ODS 出发;或保证「CDC 与批量导出」的数据范围一致(含 T+1 修正) |
| 计算口径一致 | 实时 SQL 与离线 SQL 的时间语义/去重逻辑不同 | 共享一套指标口径定义(统一 DWD 字段),实时用累积窗口对齐"当日累计",离线用同样语义;关键指标做双跑对比 |
| 数据完整性 | 实时受迟到数据、乱序影响;离线是"全量确定" | 实时按时段对账(当日实时结果 vs 次日离线结果),差异做补偿;接受"实时近似 + 离线修正" |
| 时效性 | 实时必须"当天出",离线 T+1 修正 | ADS 表用主键 upsert:离线结果可以覆盖/修正实时结果;对外口径以离线为准 |
三层保障的落地建议:
① 技术保障:实时链路端到端 EO(Source 位点 + 状态 + Sink 事务/幂等)
② 工程保障:DWD 统一 + 指标口径文档化 + 自动化对账任务(日级)
③ 业务保障:对外承诺"实时看趋势、离线看准确",避免把实时值当最终账10. Flink SQL 作业线上常见问题与排查?
答: 与 DataStream 作业的问题同源,但 SQL 有自己的几个"专属坑":
| 现象 | 常见根因 | 排查/解决 |
|---|---|---|
| 状态持续膨胀、Checkpoint 超时 | table.exec.state.ttl 没设(默认永不过期);双流 join/TopN 状态大 | 设 TTL;减小 join 范围;开 MiniBatch/RocksDB |
| 窗口不输出结果 | Watermark 不推进(Kafka 空分区、source.idle-timeout 没配) | 开 table.exec.source.idle-timeout;检查水位定义 |
| 结果少了/重复了 | 回撤消息(-U)没被下游正确处理;使用了非确定性函数 | 下游用 upsert-kafka 或支持 upsert 的表(Doris Unique) |
| 维表关联不上,大量 null | 维表数据未就绪 / 缓存了 null / 异步超时 | 缓存策略调整(未命中不缓存 null)、异步超时降级、侧输出补偿 |
| 同一个 SQL 两条流 join 结果不准 | 两条流水位不一致、慢流拖住 | 调 source.idle-timeout、给双流设置一致的 Watermark |
| SQL 作业吞吐低 | 没开 MiniBatch、没开两阶段聚合、Lookup Join 未异步 | 检查上述配置 |
| SQL 编译/计划阶段报错 | 某些 SQL 特性不支持(如 update 主键、非确定性) | 查 Flink 版本支持矩阵;改写 SQL |
| checkpoint 里有"巨量 KV" | ROW_NUMBER() 去重 / DISTINCT 状态无限增长 | 加 TTL、改成两阶段去重、控制分区数 |
定位 SQL 作业的标准动作:
① 先在 SQL 客户端(或 Flink Web SQL)里 EXPLAIN 看执行计划,确认算子与并行度
② 到 Web UI 看反压与 busyTime(与 DataStream 一致)
③ 到 Checkpoint 详情看各算子状态大小 → 找膨胀点
④ 对照上面的"SQL 专属配置"逐项检查(TTL / MiniBatch / idle-timeout / async-lookup)11. Flink SQL 与 Kafka 配合时的 Exactly-Once 怎么保证?
答:
CREATE TABLE kafka_sink (
user_id BIGINT,
total DECIMAL(10,2),
PRIMARY KEY (user_id) NOT ENFORCED -- upsert-kafka 需要主键
) WITH (
'connector' = 'upsert-kafka', -- 支持 update/delete 流
'topic' = 'user_total',
'properties.bootstrap.servers' = 'kafka:9092',
'key.format' = 'json',
'value.format' = 'json',
'sink.delivery-guarantee' = 'exactly-once', -- ★ 关键
'sink.transactional-id-prefix' = 'flink-user-total'
);要点:
| 要点 | 说明 |
|---|---|
connector = 'upsert-kafka' | 普通 kafka connector 只支持 append 流;有回撤/更新的结果必须用 upsert-kafka,它把 -U/+U/-D 编码成「key + null 值」的墓碑消息 |
sink.delivery-guarantee = exactly-once | 底层就是两阶段提交(事务预提交 + Checkpoint 完成后提交,见《Flink(三)》) |
transactional-id-prefix 必须唯一且稳定 | 多个作业用同一个前缀会 ProducerFencedException;前缀里通常带环境/作业名 |
| Source 端也要 EO | KafkaSource 的位点只在 Checkpoint 成功后才提交(不要开 auto.offset.commit) |
| 下游消费方要能处理墓碑消息 | null 值的消息代表删除;下游(如 Doris)要支持按主键删除 |
若下游是 Doris / ClickHouse / JDBC:通常不需要 Flink 事务,靠「批量写 + 主键 upsert(幂等)」实现效果上的 EO——这是更常见也更省的做法。
12. Flink 的生态组件(Source/Sink)常用哪些?
答:
| 类别 | 组件 | 说明 |
|---|---|---|
| 消息队列 | Kafka(首选)、Pulsar、RocketMQ、RabbitMQ | Kafka Source/Sink 支持 EO、位点管理、多分区并行 |
| CDC | MySQL-CDC / PostgreSQL-CDC / Oracle-CDC / MongoDB-CDC | 见第 5 题 |
| OLAP/DB | Doris(Stream Load/Connector,Unique 表 upsert)、ClickHouse、StarRocks、HBase、JDBC(MySQL/PG) | 分析型写入多用 批量 + upsert |
| KV/缓存 | Redis(常用维表点查 / 结果缓存)、HBase(大维表点查) | Redis 多用于「热点结果外置」 |
| 检索 | Elasticsearch | 日志/搜索场景 |
| 文件/对象存储 | HDFS、S3、OSS 的 FileSystem Connector,File Sink 支持滚动策略(按时间/大小滚动) | 冷数据落文件用 |
| 数据湖 | Paimon、Hudi、Iceberg | 见第 7 题 |
| 其他 | 维表用 JDBC/HBase/Redis;文本用 datagen(测试)、print(调试)、blackhole(压测) | — |
选 Sink 的三条原则:
- 能不能 upsert? 分析型(Doris/StarRocks/Paimon)用 upsert 表;仅追加场景用 Kafka;
- 怎么保证不丢/不重? 事务(Kafka EO)或幂等(主键 upsert);
- 批量 vs 延迟:批量越大吞吐越高、延迟越大;按业务 SLA 调
sink.buffer-flush.*。
13. 综合实战:「从 MySQL 到 Doris 的实时数仓」怎么搭?
答: 一道完整的架构设计题,按「链路 + 关键决策」来答:
要主动说出的四个决策(面试加分):
| 决策 | 选择 | 理由 |
|---|---|---|
| 为什么 CDC 直接入湖,不先入 Kafka? | 入湖(Paimon)长期可回溯,Kafka 只适合短期缓冲;入湖后实时与离线可共用一份 ODS | 口径统一 + 成本 |
| 为什么 DWD 用主键表? | 上游 binlog 有大量 update/delete,主键表天然 upsert,避免自己写去重逻辑 | 简单 + 正确 |
| 为什么 ADS 写 Doris 而不是直接查 Paimon? | 看板要高并发低延迟点查,Doris 更合适;Paimon 适合大范围扫描与回溯 | 场景匹配 |
| 状态怎么控制? | DWD 是无状态/短状态(去重窗口小);DWS 才是状态大户 → TTL + MiniBatch + 两阶段聚合 | 状态是成本与稳定性核心 |
Flink 系列小结:(一)架构地基 →(二)时间与窗口 →(三)状态与容错 →(四)反压与调优 →(五)SQL 与生产落地。这五篇串起来,基本覆盖了 Flink 面试「从原理到落地」的全链路,建议按顺序复习,并在每一层都能说出「为什么这么设计、代价是什么、怎么在生产中取舍」。
