引言
流处理长期被视为"高门槛":要理解状态、水位线、检查点,还要写 DataStream API 的算子链。流式 SQL(Streaming SQL)试图把这件事拉回分析师的舒适区——用 SQL 描述"持续计算",让引擎去处理状态、容错与并行。
但流式 SQL 不是把批 SQL 搬过来那么简单:批处理的表是静态快照,流处理的表是持续变化的流。理解这个差异,才能理解为什么同样的 JOIN 在流里要分三种、为什么 COUNT(*) 的输出是一条不断更新的流。本文以 Flink SQL 为主线、ksqlDB 为对照,把这两套体系讲透。
一、流式 SQL 的核心抽象:动态表
1.1 表是流的一个快照
流(Stream) = 无限的、按时间追加的事件序列
动态表(Dynamic Table) = 流在某时刻的物化视图,随新事件不断变化
stream → [持续查询] → 动态表 → 输出 changelog
关键认知:在流上执行 SQL,得到的是一个持续更新的结果表,而不是一次性结果。SELECT user_id, COUNT(*) FROM clicks GROUP BY user_id 在批里返回最终计数,在流里返回"每当某个用户点击,就输出该用户的最新计数"。
1.2 三种 changelog 模式
动态表变化通过 changelog 流表达,Flink 支持三种:
| 模式 | 含义 | 典型算子 | 能否再消费 |
|---|---|---|---|
| Append-only | 只追加,无更新删除 | 无聚合的 SELECT、窗口聚合 | 可直接当流用 |
| Upsert | 按主键 upsert/delete | GROUP BY 聚合 | 需主键,+I/-U/+U/-D |
| Retract | 先撤回旧值再发新值 | 无主键的聚合、OVER | 通用但消息量大 |
toChangelogStream 会把三种模式统一成 RowKind(+I insert、-U update_before、+U update_after、-D delete)。写 Kafka 时通常选 upsert 模式(upsert-kafka 连接器)来减少消息量:
CREATE TABLE user_click_count (
user_id BIGINT,
cnt BIGINT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'upsert-kafka',
'topic' = 'user_click_count',
'key.format' = 'json',
'value.format' = 'json'
);
二、Flink SQL 实战
2.1 DDL 与连接器
Flink SQL 把 Kafka、JDBC、Iceberg、Hive 等统一成 CREATE TABLE:
CREATE TABLE clicks (
user_id BIGINT,
page_id STRING,
click_time TIMESTAMP(3),
WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'clicks',
'properties.bootstrap.servers' = 'kafka:9092',
'properties.group.id' = 'clicks-job',
'scan.startup.mode' = 'group-offsets',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
连接器参数里的几个关键点:
| 参数 | 作用 | 常见取值 |
|---|---|---|
scan.startup.mode | 起始位点 | earliest-offset / latest-offset / group-offsets |
properties.group.id | 消费组,决定 offset 提交 | 作业名,勿复用 |
format | 消息格式 | json / avro / debezium-json / canal-json |
json.ignore-parse-errors | 脏数据容忍 | 生产建议 true + 死信侧输出 |
2.2 时间属性与水位线
流式 SQL 有三种时间:
-- 1. 处理时间(Processing Time):用引擎当前时间,无需水位线
proc_time AS PROCTIME()
-- 2. 事件时间(Event Time):从数据里取,需 WATERMARK 声明
WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND
-- 3. 事件时间列的两种写法
-- DDL 中声明(推荐,全表统一)
-- 查询中用 WATERMARK 派生
水位线(Watermark)表达"事件时间推进到哪了":click_time - INTERVAL '5' SECOND 意味着容忍 5 秒乱序,早于水位线到达的数据被视为迟到数据。水位线的容忍度与迟到策略详见 https://plumephp.com/data-streaming-window-time/。
2.3 窗口 TVF
Flink 1.13 起用窗口表值函数(Windowing TVF)替代旧的 GROUP BY TUMBLE(...) 语法:
-- 滚动窗口:每分钟每页面 PV
SELECT
window_start, window_end, page_id, COUNT(*) AS pv
FROM TABLE(TUMBLE(TABLE clicks, DESCRIPTOR(click_time), INTERVAL '1' MINUTE))
GROUP BY window_start, window_end, page_id;
-- 滑动窗口:每 30 秒统计过去 5 分钟
SELECT window_start, window_end, COUNT(*) AS pv
FROM TABLE(HOP(TABLE clicks, DESCRIPTOR(click_time), INTERVAL '30' SECOND, INTERVAL '5' MINUTE))
GROUP BY window_start, window_end;
-- 累积窗口:适合"从当天 0 点到现在的累计"报表
SELECT window_start, window_end, COUNT(*) AS cumulative_pv
FROM TABLE(CUMULATE(TABLE clicks, DESCRIPTOR(click_time), INTERVAL '10' MINUTE, INTERVAL '1' DAY))
GROUP BY window_start, window_end;
-- 会话窗口
SELECT window_start, window_end, COUNT(*) AS session_pv
FROM TABLE(SESSION(TABLE clicks, DESCRIPTOR(click_time), INTERVAL '30' MINUTE))
GROUP BY window_start, window_end;
TVF 的最大好处是窗口边界(window_start / window_end)成为可用的列,能继续参与 Join 与聚合,而旧语法只能在窗口上做单层聚合。
2.4 Join 的三种语义
流上的 Join 与批完全不同,必须理解各自的语义与代价:
| 类型 | 语义 | 状态代价 | 适用 |
|---|---|---|---|
| Regular Join | 任意一侧到达即尝试匹配,两侧状态都保留 | 无限增长 | 两侧都不大,或需完整历史 |
| Interval Join | 限定时间区间内的匹配 | 有限(按时间清理) | 订单-支付这类有时间约束的关联 |
| Temporal Join(时态表) | 按事件时间关联到当时生效的维表版本 | 维表版本保留 | 汇率、价格等随时间变化的维表 |
| Lookup Join | 每条记录查一次外部维表 | 无状态(可缓存) | 维表在 MySQL/HBase,需实时查询 |
-- Interval Join:订单与 30 分钟内的支付关联
SELECT o.order_id, o.amount, p.pay_time
FROM orders o
JOIN payments p
ON o.order_id = p.order_id
AND p.pay_time BETWEEN o.order_time AND o.order_time + INTERVAL '30' MINUTE;
-- Lookup Join:实时查 MySQL 维表
SELECT o.order_id, c.city
FROM orders AS o
JOIN dim_customer FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.user_id = c.user_id;
Regular Join 的状态无限增长是生产事故的常见来源——它需要长期保存两侧所有数据。若能用 Interval Join 或 Lookup Join 表达,务必替换。
2.5 状态与性能
流式 SQL 的状态由引擎隐式管理,但代价必须由使用者承担:
-- 开启状态 TTL,避免状态无限增长(Flink 1.18+ 表级配置)
SET 'table.exec.state.ttl' = '24h';
-- 微批聚合:用延迟换吞吐,减少对下游的写放大
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5s';
SET 'table.exec.mini-batch.size' = '5000';
-- 开启两阶段聚合,缓解热点 Key
SET 'table.optimizer.agg-phase-strategy' = 'TWO_PHASE';
参数与算子的对应关系:
| 参数 | 影响 |
|---|---|
table.exec.state.ttl | 所有有状态算子的状态保留时长 |
mini-batch | 聚合算子攒批,吞吐↑ 延迟↑ |
agg-phase-strategy | 两阶段聚合缓解热点 Key |
table.exec.sink.upsert-materialize | upsert 前是否去重 |
状态后端的选型与 Checkpoint 调优直接决定这些作业能否稳定运行,相关参数见 https://plumephp.com/data-streaming-exactly-once-state/。
三、ksqlDB 对比
3.1 架构与语法
ksqlDB 是 Confluent 推出的流式 SQL 引擎,构建在 Kafka Streams 之上,与 Kafka 深度绑定:
-- 从 topic 建流
CREATE STREAM clicks (
user_id BIGINT,
page_id VARCHAR,
click_time TIMESTAMP
) WITH (
KAFKA_TOPIC = 'clicks',
VALUE_FORMAT = 'JSON',
TIMESTAMP_FORMAT = 'yyyy-MM-dd HH:mm:ss'
);
-- 物化成表(会创建新 topic 并持续更新)
CREATE TABLE page_pv AS
SELECT page_id, COUNT(*) AS pv
FROM clicks
WINDOW TUMBLING (SIZE 1 MINUTE)
GROUP BY page_id
EMIT CHANGES;
与 Flink SQL 的差异:
| 维度 | Flink SQL | ksqlDB |
|---|---|---|
| 执行引擎 | Flink(独立集群) | Kafka Streams(Kafka 内) |
| 数据源 | 多源(Kafka/JDBC/Iceberg/Hive) | 仅 Kafka |
| Sink | 多目标 | 主要写回 Kafka topic |
| 部署形态 | 独立集群 + 作业提交 | ksqlDB Server 集群,SQL 即服务 |
| 事件时间 | 完整 WATERMARK 支持 | 基于 record timestamp |
| 生态 | 批流一体、湖仓集成 | Kafka 生态内闭环 |
3.2 pull query 与 push query
ksqlDB 把查询分成两类,这是它区别于 Flink SQL 的重要设计:
-- Push Query:持续推送结果,永不结束(流式查询)
SELECT page_id, COUNT(*) FROM clicks EMIT CHANGES;
-- Pull Query:对物化表做点查,立即返回当前值(类似数据库查询)
SELECT * FROM page_pv WHERE page_id = 'home' LIMIT 1;
Pull Query 让 ksqlDB 兼具"流处理引擎"与"实时物化视图数据库"两种角色——上游写入 Kafka,下游用 SQL 点查最新状态,无需额外的 KV 存储。代价是物化表必须常驻内存/磁盘,规模受限于集群资源。
四、选型对照
| 场景 | 推荐 | 理由 |
|---|---|---|
| 多源入湖 + 批流一体 | Flink SQL | 连接器丰富,Iceberg/Hive 原生支持 |
| 纯 Kafka 内的实时转换与聚合 | ksqlDB | 部署轻、SQL 即服务、无需独立集群 |
| 需要实时点查最新状态 | ksqlDB | Pull Query 开箱即用 |
| 复杂事件处理 / 自定义算子 | Flink DataStream | SQL 表达力不足时下探 API |
| 已有 Flink 平台与运维体系 | Flink SQL | 复用平台能力与监控 |
一条务实建议:若技术栈已围绕 Kafka 且需求是"实时清洗 + 聚合 + 写回",ksqlDB 的上手成本更低;一旦涉及湖仓写入、多源关联或批流统一,Flink SQL 是唯一解。
五、生产实践
5.1 从 CDC 到实时指标
一条典型的实时链路:CDC 采集 → 流式 SQL 聚合 → 写入 OLAP。
-- 1. 消费 Debezium 变更流
CREATE TABLE orders_cdc (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(12,2),
status STRING,
ts TIMESTAMP(3),
WATERMARK FOR ts AS ts - INTERVAL '10' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'dbserver1.inventory.orders',
'format' = 'debezium-json',
'properties.bootstrap.servers' = 'kafka:9092'
);
-- 2. 每分钟 GMV(只统计已支付)
INSERT INTO gmv_per_minute
SELECT
window_start, window_end,
SUM(amount) AS gmv,
COUNT(DISTINCT user_id) AS buyers
FROM TABLE(TUMBLE(TABLE orders_cdc, DESCRIPTOR(ts), INTERVAL '1' MINUTE))
WHERE status = 'paid'
GROUP BY window_start, window_end;
CDC 的格式解析与 Schema 演进细节见 https://plumephp.com/kafka-connect-cdc/;这条链路的存储层与整体架构见 https://plumephp.com/realtime-data-warehouse/。
5.2 部署与运维要点
| 要点 | 做法 |
|---|---|
| 并行度 | 按 Kafka 分区数设置,避免空转或消费不足 |
| 资源隔离 | 每个作业独立集群或独立 slot,避免相互影响 |
| 状态大小 | 优先用 TTL + 微批控制,避免 Regular Join 无限增长 |
| 结果一致性 | 开 Checkpoint 并配 upsert-kafka 保证端到端精确一次 |
| 可观测性 | 监控 source lag、checkpoint 时长、状态大小 |
5.3 SQL 作业的版本化与测试
SQL 作业也是代码,同样需要版本管理与测试。Flink 生态可用的手段有三类:
1. 语法与计划校验 提交前用 sql-gateway / TableEnvironment 编译一次,拦住语法错
2. 确定性输入输出测试 用 datagen 连接器或固定输入表,跑一遍断言 changelog
3. 结果 diff 影子作业双跑,对比输出 topic 的 changelog 是否一致
用 datagen 构造确定性输入做回归测试:
-- 测试用输入:固定 10 条,可复现
CREATE TABLE clicks_test (
user_id BIGINT, page_id STRING, click_time TIMESTAMP(3),
WATERMARK FOR click_time AS click_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'datagen',
'number-of-rows' = '10',
'fields.user_id.kind' = 'sequence',
'fields.user_id.start' = '1',
'fields.user_id.end' = '10'
);
把作业的 SQL 拆成"源表定义 + 查询逻辑 + 目标表定义"三段并纳入 Git,配合 CI 编译校验;上线走"影子作业 → changelog diff → 切换",与批处理管道的发布纪律一致。
六、踩坑清单
| 坑 | 表现 | 修法 |
|---|---|---|
| Regular Join 状态爆炸 | Checkpoint 越来越大直到失败 | 换 Interval Join 或 Lookup Join |
| 忘设水位线 | 只能用处理时间,窗口结果不可复现 | DDL 里声明 WATERMARK |
| 处理时间当事件时间 | 迟到/重放数据指标错乱 | 业务指标一律用事件时间 |
group.id 复用 | 多作业争抢位点,数据丢失 | 每作业唯一 group.id |
| 脏数据导致作业失败 | 单条坏消息让作业反复重启 | ignore-parse-errors + 死信队列 |
| 无 TTL 的状态 | 长跑后状态失控 | table.exec.state.ttl |
| 维表变更不生效 | Lookup Join 缓存过期时间长 | 调小 cache TTL 或加失效通知 |
小结
流式 SQL 的价值在于把"持续计算"用声明式语法表达,让实时开发从写算子变成写查询。掌握它的关键是三件事:理解动态表与 changelog 的语义(结果是流而非快照)、理解流上 Join 的三种语义与状态代价(Regular Join 慎用)、理解时间属性与水位线(业务指标必须用事件时间)。Flink SQL 胜在连接器生态与批流一体,ksqlDB 胜在轻量与 Kafka 内闭环,选型取决于数据源与目标是否超出 Kafka 边界。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。