流式 SQL:Flink SQL 与 ksqlDB

流式 SQL 把无限流当作持续变化的表,用标准 SQL 描述持续计算,是降低实时开发门槛的关键。本文讲清动态表与 changelog 的三种模式、Flink SQL 的 DDL 与连接器定义、事件时间与水位线声明、窗口 TVF 与各类 Join 的语义差异、状态与性能代价,再对比 ksqlDB 的架构、语法与 pull/push 查询,给出选型对照与生产踩坑清单。

引言

流处理长期被视为"高门槛":要理解状态、水位线、检查点,还要写 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/deleteGROUP 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'
);

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-materializeupsert 前是否去重

状态后端的选型与 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 SQLksqlDB
执行引擎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 即服务、无需独立集群
需要实时点查最新状态ksqlDBPull Query 开箱即用
复杂事件处理 / 自定义算子Flink DataStreamSQL 表达力不足时下探 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 边界。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 湖仓访问控制与权限治理
  2. 非结构化文档 ETL 与多模态数据
  3. Flink 状态后端与 Checkpoint 调优