我理解的流批一体:从三条链路到一个参数

这篇文章由和claude反复对话生成。全文用同一张订单宽表贯穿,从现状的两条链路讲到 Fluss + Paimon 的改造,最后落到 Materialized Table 的增量语义。

先给出立场

"流批一体"这个词至少被用在三个不同层次上,讨论之前得先说清楚在讲哪一层:

层次

主张

代表

引擎层

一套引擎既能跑流又能跑批

Flink 流批统一的 API 与运行时

存储层

一份存储既能流式读写又能批读写

Paimon / 湖仓一体

语义层

同一份业务语义,在不同时效要求下保持一致

本文的主题

前两层是手段,第三层才是目的。而语义之所以统一不了,根子在开发成本:实时开发的单位成本远高于离线,高到我们不得不为它单独设计一套简化的逻辑——同一张表于是长出了两条链路、两套代码、两种口径。

后面三个部分,就是把这笔成本逐层剥开:

  • 第一部分盘现状:两条链路各自都有充分理由,但它们写同一张表,于是有了重复实现、口径分叉和写入冲突;而每一种缓解冲突的办法,都是再叠一层人工约束。

  • 第二部分用 Fluss + Paimon 把两条链路合成一条:实时开发的成本降下来,实时链路不必再为了活下去而简化,口径分叉的根源随之消失;写入者收敛回一个,冲突在物理上不再成立。

  • 第三部分用 Materialized Table 统一语义:一份定义加一个新鲜度声明,执行形态交给框架决定——写这张表的人不需要知道它会被跑成流式还是批式。

第一部分:现状 —— 两条链路是怎么长出来的

讲流批一体最容易犯的错误,是一上来就摆一张画满箭头的架构图,然后说"你看,Lambda 架构就是有问题"。这个讲法有个隐含前提:现在的架构是设计失误。

但真实情况是,两条链路不是一次设计出来的,而是被两个各自独立、各自成立的诉求逼出来的。先讲清楚每条链路当初为什么必须存在,后面的改造才有意义。

一、一个例子

StarRocks 上的订单宽表 dws_order_wide,主键 order_id一张当前状态表——每个订单只保留最新状态,不按天分区。它服务客服工作台、运营大屏和即席查询。

字段就是订单的那些常规字段:order_statuspay_amtcreate_timecate_id 等等。

离线链路的加工结果落在 Hive 上,再同步到 StarRocks 对外服务;实时链路则从 Kafka 消费事件后直接写 StarRocks。两条链路是两套独立的代码,也是这张表的两个写入者。

需要回看历史状态时(同比、回溯、对账),只能另外每天存一份全量快照分区,存储按天线性增长。

二、两条链路各自的合理性

实时链路:变更必须马上可见。 订单状态变更要秒级出现在客服工作台。用户打电话说"我付款了状态怎么还没变",客服看到的必须是当前态。天级满足不了,小时级也满足不了。这是业务硬约束,不是选型偏好。

但它的逻辑是刻意简化过的:只读事件流,再点查 HBase 把维度打宽,不做多流 join。核心矛盾在双流 join——两边都要存全量 state,状态规模压不住;而且任何一边的更新都会触发回撤,而回撤是"先撤旧值、再发新值"两条消息,这中间的一瞬间下游看到的就是错的:行短暂消失,或者聚合值先掉下去再涨回来。查询随时都在发生,正好落在这个窗口上就会读到错值,事后还查不到、复现不了。为了让这条链路能被维护,我们在逻辑上做了一些简化。

离线链路:需要一份会收敛的全量。 实时链路的正确性依赖太多外部条件——Kafka 可能丢消息、作业可能挂掉丢状态、CDC 可能漏采、上线可能引入 bug。而这些问题在实时链路里没有自愈机制:一行写错了就一直错下去。所以必须有一条从上游全量快照出发、可重跑、可校验的链路,每天把整表覆盖一遍。它的价值不是时效,是收敛性:无论昨天发生了什么,第二天早上数据一定是对的。这条链路用 Spark SQL 做分层加工,跑的是天级批任务。

两个诉求:低延迟、正确性。彼此独立,各自成立。

(现实里往往还有第三条小时级链路,用来刷新那些"T+1 太慢、但上游只能到小时"的指标。它本质上是离线链路换了个调度周期,下面所有的论证对它同样适用,为了聚焦这里就不单列了。)

flowchart LR DB[("MySQL")] -->|CDC| MQ[("Kafka")] MQ --> F1["Flink SQL 实时链路<br/>事件流 + 点查打宽"] DIM[("HBase<br/>维表")] -.->|lookup| F1 MQ --> ODS["Hive ODS"] ODS -->|Spark SQL 天级| DWS["Hive DWD / DWS"] F1 -->|"整行 upsert · 秒级"| SR DWS -->|"全字段 · 整表覆盖 · 天级"| SR SR[("StarRocks<br/>dws_order_wide<br/>当前状态表")] --> APP["客服工作台<br/>运营大屏"] classDef rt fill:#e3f2fd,stroke:#1976d2,color:#000 classDef off fill:#fff3e0,stroke:#f57c00,color:#000 class F1 rt class ODS,DWS off

三、然后问题出现了

两条链路各自合理,写同一张表同样合理——使用方要的就是一张能直接查的宽表,不该也没法感知某个字段从哪条链路来。所以问题不在"为什么写同一张表",而在于:两个写入者、两种写入语义、两种时效,落在同一张表上,却没有任何机制协调它们。

1. 两套逻辑、六个组件,开发成本高

最直接的成本是同一份业务语义写了两遍:Flink SQL 一遍、Spark SQL 一遍,两种引擎方言。改一次口径要改两处,漏改一处就是长期存在、难以发现的不一致。

但更麻烦的是这两遍并不等价。离线链路能做复杂的多表 join,实时链路只能点查打宽——所以有些字段在实时侧算起来代价太高,索性交给离线补;有些字段两边都算得出,但走的是完全不同的路径。两边的差异不是 bug,是设计使然。 也正因如此,两条链路的结果永远对不齐——你没法要求两套本来就不同的逻辑给出一样的结果。

组件账也不轻:Kafka、HBase、Flink、Spark、Hive、StarRocks——一个订单宽表的字段,要经过六个系统。每加一个字段,都要判断它归哪条链路、会不会和另一条冲突、维表要不要同步扩一列。

而两条链路里,实时那条的单位成本远高于离线:

离线开发

实时开发

时间语义

只有 dt,跑批时数据已就位

事件时间 / 处理时间 / watermark,还要处理迟到

状态

无状态,重跑即可

TTL、大小、backend 选型、能否从旧 state 恢复

Join

写就完了

双流 join 的 state 压不压得住?维表 as-of 对不对?

回撤

不存在

join 和无界聚合必然产生,中间态会不会正好被查到

上线

改 SQL 重跑

改 SQL → state 不兼容 → 怎么迁移

验证

跑一次看结果

只能等数据来,或者自己造事件流

出错排查

重跑,中间层随便查

乱序、迟到、state 过期还是维表版本?而中间结果锁在 state 里查不到

离线开发是写一段 SQL;实时开发是写一段 SQL,再回答七个关于执行的问题。 所以同一个人做实时的产出可能只有做离线的三分之一,而能做实时的人本来就少。

这就是上面那个"刻意简化"的由来——不是能力选择,是成本倒逼。而简化的代价,就是口径先天分叉。

2. 全量覆盖的丢失窗口

上面是成本问题,下面是正确性问题。

离线任务凌晨 01:00 读取上游全量快照,此刻订单 O1 是"未支付";跑批、加工、同步走了两小时;02:00 用户付款,实时链路已经把 StarRocks 写成"已支付";03:00 离线结果整表覆盖,把它退回"未支付"。

sequenceDiagram participant K as Kafka participant RT as 实时作业 participant D as StarRocks 宽表 participant BATCH as 离线链路 Note over D: 订单 O1 当前为「未支付」 BATCH->>BATCH: 01:00 读取上游全量快照 rect rgb(255, 240, 240) K->>RT: 02:00 支付事件 RT->>D: upsert,状态变为「已支付」 BATCH->>D: 03:00 跑批完成,整表覆盖 Note over D: 退回「未支付」<br/>直到该行下次变更才修正 end rect rgb(240, 250, 240) BATCH-->>RT: 覆盖成功后触发重启 RT->>K: offset 重置到读取点之前 K->>RT: 重放窗口内事件 RT->>D: upsert,恢复为「已支付」 end

覆盖是整表的,所以影响范围不是某一类订单,而是读取快照到覆盖完成这段窗口内所有发生过变更的行。这些行会停在旧值上,直到它下一次发生业务变更才被实时链路带回正确状态——如果它已经进入终态、之后长期不变,那就是长期错误。

凌晨订单量低,绝对数量看着不大,而且这些行"迟早会自己好",所以它长期不被当成一个问题——直到某天客服反馈某个订单状态一直不对。

根因不是覆盖本身,而是胜负由到达顺序决定,而不是由数据版本决定。这是典型的 last-write-wins by arrival time:只要存在多个写入者、且各自数据版本时刻不同,就必然丢数据。

3. 每个解法,都是再加一层人工维护的约束

丢失窗口有两条解法:加版本列,或者重放事件

解法一:版本列。 给结果表指定一个 sequence 字段——取业务行自身的 update_time——让引擎按版本大小定胜负,而不是按到达顺序。离线覆盖携带的是快照时刻的旧版本,天然写不进那些已经被流式更新过的行。

flowchart TB A["离线整表覆盖<br/>携带 01:00 的 update_time"] --> C{"引擎比较版本列"} B["实时 upsert<br/>携带 02:00 的 update_time"] --> C C -->|"版本更大 → 写入"| D[("结果表<br/>已支付")] C -->|"版本更小 → 丢弃"| E["旧值被挡住"] classDef off fill:#fff3e0,stroke:#f57c00,color:#000 classDef rt fill:#e3f2fd,stroke:#1976d2,color:#000 class A off class B rt

机制很简单,但它有三个前提,每一个都不由数仓控制:

一是存储得支持。 StarRocks 的主键表可以指定 sequence 列,MySQL 这类只认"后写覆盖先写"的存储就没有这个能力——换个 sink,这条路直接走不通。

二是版本字段没得选。 处理时间等于没做,又回到按到达顺序;CDC 的位点或采集时间离线侧拿不到,两个写入者之间不可比。唯一两边都能拿到、且语义同源的就是业务行自身的 update_time

三是这个字段的可靠性由上游决定。 触发器直接改表、DBA 手工修数据、ORM 的选择性更新——任何一条不更新 update_time 的写入路径,都会让版本比较失效。而且是静默失效:不报错,只是某些行的值一直不对。

还有一个次要代价:写入从覆盖退化成 merge,行集合不再严格对齐。上游物理删除、而流式恰好漏掉 -D 的那些行会永远留着——覆盖时它们本来会被顺手清掉。

解法二:重放事件。 不加版本列,而是在离线覆盖完成后触发实时作业重启,把 offset 拨回快照读取点之前,让窗口内的事件重新走一遍。因为写入者始终只有一个、事件按顺序依次重放,最终值必然正确。

问题出在"最终"这两个字上——catch-up 期间表处于一个对外可见的错误中间态:它要按顺序把几小时前的事件重新放一遍,这段时间查询会看到数据先退回旧值、再逐步恢复,而且不同行的恢复进度还不一样。所以重放不能直接在对外的那张表上做,需要配一套双链路切换:覆盖和重放都在备表上完成,追平并校验通过后再把读流量切过去。

sequenceDiagram participant Q as 查询方 participant A as 表 A(当前对外) participant B as 表 B(备用) participant RT as 实时作业 participant BATCH as 离线作业 Note over Q,A: 稳态:查询读 A,实时写 A RT->>A: 持续 upsert rect rgb(255, 248, 240) Note over BATCH,B: 01:00 离线读取快照,开始跑批 BATCH->>B: 03:00 整表覆盖写入 B Note over B: B 此刻缺失 01:00 之后的变更 end rect rgb(232, 244, 253) RT->>B: 从 01:00 之前的 offset 重放,追平 B Note over B: 追平期间 B 不对外,<br/>中间态无人可见 RT->>A: 同时继续写 A,查询不受影响 end rect rgb(240, 250, 240) Note over B: 追平完成,校验通过 Q-->>B: 切换读流量到 B Note over A: A 转为备用,下一轮角色对调 end

代价是双份存储、一层切换编排,以及对作业的硬要求:必须无状态,或状态可从重放中重建——一旦有跨天累计、全局去重、双流 join,重置 offset 就是丢 state,结果直接算错。

两条路都有效,很多团队就这么跑了好几年。但把它们放在一起看,是同一个模式:每一次缓解,都是在架构之上再叠一层需要人工维护的约束。 update_time 必须被上游可靠维护、重放起点必须早于快照读取点、切换前后两张表的保留期要对齐、有状态的作业不能用重放——这些约束没有一条写进任何系统里,只存在于文档、注释和某个人的记忆中。正确性依赖的是"你记得所有约束",而不是"系统保证了它们"。

4. 时效性被硬编码进了架构

现在这条边界是显式可见的:秒级的走 Kafka,天级的走 Spark。 哪个字段要多快,决定的不是一个参数,而是它归哪条链路、用哪套代码、走哪套调度。

如果哪天某个指标要从天级提到分钟级,或者反过来实时链路成本太高要降到五分钟——答案都是把字段在链路之间搬家,重写、重配、重新验证。

时效性本该是一个参数,现在却是一个架构决策。 这不是某个组件的能力缺陷,是整套架构的表达能力缺陷。

四、小结

两条链路各有充分理由,已有的工程实践也确实压住了大部分问题。现状不是不能用。

但两件根本的事从头到尾没被解决:

  • 开发成本:同一份语义写两遍、跨六个组件,而且实时那遍因为成本太高被迫简化,导致口径先天分叉、差异不可解释。

  • 协调成本:两个写入者落在同一张表上,冲突全靠人来编排,而这些编排只存在于文档和某个人的记忆里。

这两条其实是同一件事的两面:如果实时链路不必再简化,它就不需要离线那条来补齐字段;如果只剩一个写入者,冲突也就无从谈起。 所以下一步的着力点是让实时开发不再那么贵——链路合成一条是这件事的结果,不是先定下的目标。

第二部分:Fluss + Paimon —— 把协调下沉进存储

两条链路各有各的诉求,硬合只是把问题挪个地方。所以先看清每个诉求该由谁承接

一、职责重新划分

诉求

现状由谁承担

改造后由谁承担

秒级可见

独立的 Kafka 链路

Fluss 日志层(同一张表的热层)

收敛性兜底

Spark 天级整表覆盖

从表重建 + 版本比较

历史回溯

每天另存全量快照分区

Paimon tag 时间旅行

关键在第一行:秒级不再需要一条独立链路,它变成同一张表的一个存储层。 这是整个阶段二的支点,后面所有收益都从这一条推出来。

二、为什么需要 Fluss

改造的第一反应通常是:把两条链路都建到湖上,用 Flink 流式写 Paimon,一套 SQL 两种执行模式,问题不就解决了?

不行,因为 Paimon 到不了秒级,而且这不是调参能解决的。

Paimon 的可见性单位是 snapshot,而 snapshot 由 checkpoint 触发提交。 数据写进去不算数,必须等 checkpoint 完成、commit 算子写完元数据,读者才能看到。端到端延迟的下限就是 checkpoint interval + commit 耗时——这是一致性模型定义的,不是实现瑕疵。

checkpoint 压不到秒级,是四个约束同时起作用:

  • 小文件产生速率:每个 checkpoint、每个 bucket 至少落一个文件。interval 从 1 分钟压到 5 秒,文件产生速率翻十几倍,compaction 追不上,读放大恶化;而 compaction 又要抢资源,形成正反馈。

  • 元数据膨胀:每次 commit 写 manifest,snapshot 数线性增长,快照过期和 manifest 合并本身开始成为负担。

  • 对象存储的提交延迟:S3/OSS 上一次 commit 涉及多轮元数据读写,本身就是百毫秒到秒级。

  • 主键表的写放大:LSM 要排序合并,小批量高频写入恰好是它最不擅长的模式。

所以工程上的舒适区是 1~5 分钟。这是"周期性提交不可变快照"这个模型的固有下限。

所以湖之上还需要一个日志层来承接秒级——写入即可见,不依赖 checkpoint 做可见性。现状里这个位置站着 Kafka,那为什么要换成 Fluss?

"秒级"本身不是理由。 Kafka 也能做到秒级,现状那条链路就是靠它撑起来的。如果只是要低延迟,换成 Fluss 没有意义。

真正的差别是一句话:Kafka 是链路外的一个独立系统,Fluss 是同一张表的一个存储层。 落到四个具体能力上:

① 热冷同源,查询侧只有一个数据源。 Fluss 的 tiering service 按 offset 顺序把日志单向归档到 Paimon,两者是同一张表的两个时效档位。用 Kafka 的话,你需要另写一个作业把数据落湖——于是湖和 Kafka 是两个写入路径,查询侧要么选一个、要么自己拼,第一部分那些多写入者问题会原样重现。

② 列裁剪的投影下推。 Fluss 是列存日志,下游只读三个字段就只传三个字段。Kafka 是整行传输,读一列也要传整行再丢掉。当一张 DWS 宽表的 changelog 被多个下游作业消费、每个只关心不同的几列时,这个差异是数量级的。

③ 主键点查,于是维表和 Delta Join。 Fluss 主键表支持毫秒级点查,带来两个直接后果:维表不再需要单独放 HBase——Fluss 表自己就是维表,和事实流是同一份数据、同一个口径;双流 join 可以退化成"来一条查一条",两边的 state 都不用存。大状态是实时数仓最常见的成本黑洞,这条的价值往往超过前两条。

④ changelog 语义原生完整。 Paimon 生成 changelog 需要 lookup 或 full-compaction changelog producer,有额外开销和延迟;Fluss 的日志天生就是 changelog,retract 语义原生完整。做多层流式 DWS 聚合时这个差别很实在。

反过来也要讲清楚边界在哪:只要有秒级诉求,日志层就是必需的——这是上面那条分钟级下限的直接推论,不存在"用 Paimon 调调参数也能凑合"的选项。但如果你的场景分钟级就够,那么只做 Paimon 那半边完全成立,上面四条能力换不来一个新组件的运维成本。我们引入它,是因为客服工作台的秒级诉求是硬的,而它一旦引入,实时侧就顺势建成了真正的分层数仓,而不只是把一条流写进湖。

三、架构

flowchart LR DB[("MySQL")] -->|CDC 增量| ODS["Fluss ODS"] SNAP[("上游全量快照")] -->|定期补数| ODS ODS --> DWD["Fluss DWD"] --> DWS["Fluss DWS<br/>主键表 · sequence"] DWS -.->|"tiering<br/>按 offset 单向归档"| PM["Paimon 归档层"] PM -->|每日打 tag| TAG[("历史快照<br/>tag 时间旅行")] DWS --> SR["StarRocks<br/>Fluss Catalog"] PM --> SR SR -->|"union read · 秒级"| APP["客服工作台<br/>运营大屏"] TAG -->|"Paimon-only · 分钟级"| HIST["回溯 / 同比 / 对账"] classDef hot fill:#e3f2fd,stroke:#1976d2,color:#000 classDef cold fill:#fff3e0,stroke:#f57c00,color:#000 class ODS,DWD,DWS hot class PM,TAG cold

四个变化点:

分层第一次同时存在于流和批。 ODS/DWD/DWS 建在 Fluss 主键表上,一套 Flink SQL 写完就是秒级的;同一批表通过 tiering 自动落 Paimon,历史数据可批读、可被其他团队接走。现状里"实时侧没有分层、中间结果锁在 Flink state 里"这条消失了。

StarRocks 不再持有物化副本。 通过 Fluss Catalog 直接 union read,一次查询同时读日志和湖。没有副本,就没有同步,就没有覆盖,也就没有第一部分那个全量覆盖的丢失窗口。 数据只有一份,新鲜度由 tiering 边界决定,而不是由某个同步任务的调度时刻决定。

两条读路径,各司其职。 当前状态走 union read(秒级),历史回溯走 tag(Paimon-only 视图,分钟级)。历史查询本来就不需要秒级,所以这两条路径天然解耦。

我们的场景 QPS 不高,union read 这条路径够用;高并发点查场景仍然需要在 StarRocks 内表留一份热区,最终形态会是混合的。

社区进度注记:StarRocks 的 Fluss Catalog 列在 2026 roadmap(starrocks#67606);Fluss 侧的 union read 已集成 Flink,面向 StarRocks/Spark/Trino 的原生 union read 同样在 2026 roadmap 上。

四、三个机制

1. union read 怎么保证不重不漏

这是最容易被怀疑的一点——两个存储层各读一半,边界怎么对齐?

答案是边界不是时间,而是 offset。tiering service 按日志顺序单向归档,提交时除了写 Paimon snapshot,还把每个 bucket 的最后归档 offset 上报给 Fluss coordinator;客户端查询时从 coordinator 拿到确切边界。湖那一半读到该位点为止,日志那一半从该位点之后开始读,既不重叠也没有空隙。

这和"按业务字段切分视图"有本质区别:业务字段会漂移(时间字段缺失时回退处理时间,同一行在两侧对不上,union 出重复行);offset 是日志本身的序号,不依赖任何业务语义,也不可能漂移。

流式读同理:先从湖里追历史(批量、高吞吐),追到边界后无缝切到日志继续消费。第一部分里"下一个需求只能从 Kafka 重新消费一遍"的问题,在这里变成一次高效的 catch-up read。

2. 兜底:从表重建

收敛性这个诉求不能丢。流式链路依然可能丢消息、丢状态、上线 bug,依然需要一条能把数据拉回正确的路径。

关键前提是:每一层都是一张 changelog 完整、可以从头读的表。所以修复的粒度是"层",不是"整条链路":

  • 某层算错了——SQL 有 bug、口径变了、作业挂掉丢了 state。改完之后让这一层从它的上游表从头消费一遍就行,重建完成,上下游都不用动。

  • 源头缺数据——CDC 漏采、上游 binlog 被清理。这部分日志里本来就没有,从表里读不出来,只能拿上游快照补进 ODS,然后它自然往下流。

两种修复都不用停掉正在跑的流式作业,因为版本列决定了它们不会互相覆盖:版本更新的行留下,更旧的丢弃。胜负由数据版本决定,不由到达顺序决定。 版本列只加在会被补数写到的那一层,中间层只有一个写入者,不需要。

这里有个关键细节:修复一个算错的值时,重建读的是同一份上游数据,update_time 不变,所以它携带的版本和表里那行是相等的,不是更大的。能不能覆盖,取决于引擎在等值时的行为。Fluss 的 versioned merge engine 是新行版本大于或等于旧行时新行获胜,所以等值的重建能盖掉错误值。三种情况刚好各归其位:

场景

版本比较

结果

漏了整行

表里没有

插入

值算错了

相等

重建覆盖

数据比现有的旧

更小

丢弃

不过有两件事这套机制盖不住。

一是删除。 versioned merge engine 不支持 DELETE,默认把删除消息忽略掉。所以上游物理删除的行,在这张表里删不掉;而快照补数时,快照里本来就没有这行,补数对它同样是 no-op。删除既不会被执行,也不会被修正。 最干净的化解是推动上游改用逻辑删除——删除退化成一次普通更新,走正常的版本比较,特殊性就消失了。

二是重建期间的中间状态。 重建不是原子的,它按顺序把行一条条写回去,不同行的进度不一样。补数场景影响小——同版本覆盖,值本来就该一样;但口径变更那种重建,表里会有一段时间新旧口径混着,跨行的聚合查询结果不可用。这类动作要避免暴露,只能在另一张表上重建完再切读流量,粒度粗,好在它本来就低频。

3. 历史回溯:用 tag 替代快照分区

现状要保留历史,只能每天存一份全量分区,存储按天线性增长。

Paimon tag 是增量存储、全量语义:tag 之间共享底层 LSM 文件,没变化的数据只存一份。一年 365 个 tag 的存储成本远低于 365 份全量分区,而查询语义和"那天的全量快照"完全一致。

这里有一个必须讲清楚的区分——在 Fluss tiered 表上,打 tag 安全,fast-forward 不安全

操作

在 tiered 表上

原因

打 tag / 读 tag

安全

只读的元数据标记,不动 main 指针

branch + fast-forward

不安全

改 main 指针,coordinator 记录的 offset 边界失效

外部 INSERT OVERWRITE

不安全

引入第二个写入者

原因在于 union read 的边界权威副本在 coordinator 手里,而不在 Paimon 里。你在 Paimon 侧把 main 指向另一个 snapshot,coordinator 并不知情,union read 会从错误的位点去接日志——要么重复、要么丢一段,而且不报错。此外 Fluss 给自己产生的 snapshot 打了专属的 commit-user 和 offset 元数据,Paimon 表还比 Fluss 表多出 __bucket / __offset / __timestamp 三个系统列,外部批作业产出的 snapshot 无法正确携带这些信息。

结论:Fluss tiered 的 Paimon 表是 tiering service 独占的归档层,只能读、只能打 tag,不能改写。 这也解释了为什么修正只能从上游重新写入 Fluss,而不是去改湖上那份归档。

五、同一个订单,在新架构下

回到第一部分那个例子:订单 O1 在离线快照读取之后、覆盖完成之前发生了支付。

时刻

发生了什么

查询看到的状态

01:00

上游快照生成,此刻 O1 = 未支付

未支付

02:00

支付事件 append 进 Fluss 日志,offset 更大

已支付(秒级可见)

03:00

快照补数写入,携带 01:00 的旧 update_time

已支付(旧版本写不进去)

任意时刻

tiering 把这段日志归档进 Paimon

已支付

没有任何一个时刻会退回旧值。

第一部分那张丢失窗口的时序图,在这里没有对应物。这是整个改造最值得记住的一句话:问题不是被修好了,是在物理上不成立了。

六、逐条对账第一部分的问题

现状问题

改造后

靠什么

三套逻辑、六个组件,开发成本高

大幅缓解

一套 SQL 两种执行模式;组件减到三个

实时被迫简化,口径先天分叉

解决

维表就是表、Delta Join 免大状态,实时可做复杂 join

实时开发本身的门槛

没解决

见第七节

全量覆盖的丢失窗口

消失

查询侧无副本可覆盖;修复靠版本比较而非覆盖

层层人工约束

大幅缓解

双链路切换不再需要;水位收敛成 offset 和 tag

历史回溯要额外存全量

解决

tag 增量存储

时效性硬编码

没解决

见第七节

前几条都出自同一个原因:协调机制从人手里下沉到了存储层。 现状里那些约束——update_time 要被上游可靠维护、重放起点要早于快照读取点、切换前后两张表的保留期要对齐、有状态的作业不能重放——要么消失,要么变成表属性和调度依赖这种可以被系统检查的东西。

七、必须诚实讲的三件事

时效性依然是硬编码的。 你还是要决定哪张表建在 Fluss、哪张只落 Paimon、哪张走批模式刷新。算法侧哪天能做到分钟级,你还是要改 DDL、改作业、改调度。换了地基,但表达能力的缺口在原地——这正是第三部分的入口。

tag 的时间语义是近似的。 tag 标记的是某个 snapshot,而 snapshot 是 tiering 的提交批次边界。所以"某天"这个 tag 的真实含义是"归档到某个 offset 时的状态",误差量级是 tiering freshness 加数据到达延迟。对天级快照完全够用,但别对外宣称它是精确时点。另外打 tag 必须排在补数之后(调度依赖),tag 会 pin 住文件,保留策略要设——常见做法是日 tag 留一年、月末 tag 留三年。

若干工程约束没有消失。 lookup join 的 as-of 语义问题一点没变,重建时维表已是新版本,结果不可复现;union read 走 merge-on-read,需要 deletion vectors 和排序优化配合;Fluss 是新组件,要单独运维、监控、排障。

八、一句话结论

现状是"两条链路各自正确,靠人来协调";这一步是"一条链路,协调下沉进存储"。

整套设计最终收敛成三件事,各自只有一个机制:秒级当前状态靠 union read,历史回溯靠 tag,数据和逻辑的修正都靠从表重建加版本比较。没有冷热分离,没有视图水位,没有影子表门禁。

开发成本这一侧也真实下降了:组件从六个减到三个,维表不再是单独的 HBase 而就是表本身,Delta Join 免掉了大状态,中间结果第一次可查可复用——实时链路不再需要为了活下去而刻意简化,它可以做和离线一样复杂的 join。口径先天分叉这个问题,到这里才算真正消掉。

但有一件事这一步动不了:写这张表的人,仍然得是懂流式执行的那一个。

第三部分:增量语义 —— 一点展望

回头看第一部分那张实时 vs 离线的对照表,七个问题——时间语义、状态、join 策略、retract、上线迁移、验证方式、排查路径——第二部分只解决了其中两三个:维表就是表、Delta Join 免掉大状态、中间结果可查。剩下的原封不动,而它们才是决定"一个需求要排多久"的那部分。

所以实时需求仍然只能由会写实时 SQL 的人来做,团队产能没有变化。这是整笔成本账里最贵的一块,也是这一节要谈的东西。

这一节是展望,不是方案——相关能力还在演进中,这里讲的是方向。

一、第二部分留下了什么

除了实时开发这条,其余的约束也只是换了形态:

约束

现在存在于哪里

谁来保证

哪张表建在 Fluss、哪张只落 Paimon

建表 DDL

人的判断

快照补数多久跑一次

调度配置

人的判断

打 tag 必须排在补数之后

调度依赖

人配的依赖

tag 保留多久

表属性 + 清理任务

人的判断

sequence 字段是否可靠

上游的写入习惯

没人保证

这六条比第一部分那些好得多——它们至少是显式配置,可以被 review、被监控。但它们仍然是分散的、彼此独立的决策,而且没有一处表达了它们真正想说的东西

比如"快照补数每天跑一次",它真正想说的是"这张表的数据可以容忍一天的收敛延迟";"这张表建在 Fluss",真正想说的是"这张表要秒级新鲜"。我们一直在用实现手段表达意图,而不是直接声明意图。

这两件事其实是同一件事的两面:你要么懂执行细节,要么做不了实时需求——不管这个细节是 watermark 还是补数节奏。

二、Materialized Table:把刷新方式从表定义里分离出来

Flink 的 Materialized Table 把这件事翻转过来:你只声明结果应该是什么、以及要多新,剩下的交给框架。

一张表的定义包含两部分——一段描述业务语义的查询,和一个 FRESHNESS 声明。可以先这样锚定:它像数据库里的物化视图,但多了一个"要多新"的维度。物化视图只回答"这张表是什么",Materialized Table 还回答"它需要多新",而后者决定了它怎么被算出来。

放到我们这张订单宽表上,写出来是这样:

CREATE MATERIALIZED TABLE dws_order_wide (
  PRIMARY KEY (order_id) NOT ENFORCED
)
FRESHNESS = INTERVAL '5' SECOND
AS SELECT
  o.order_id, o.user_id, o.order_status, o.create_time,
  p.pay_amt, p.pay_time
FROM ods_order AS o
LEFT JOIN ods_payment AS p
  ON o.order_id = p.order_id;

这里是一个双流 join——第一部分里实时链路做不起、只好交给离线的那种。写在这里就是一行 LEFT JOIN,没有 state TTL、没有 watermark、没有 sink 配置、没有作业提交,只有一段 SQL 和一个"要多新"。哪天大屏说分钟级够用了,改的就是那个 INTERVALREFRESH_MODE 也可以显式写成 CONTINUOUSFULL 来覆盖框架的判断。

框架根据 freshness 决定执行形态:要求秒级或分钟级,跑成一个持续的流式作业做增量刷新;要求小时或天级,退化成周期性的批作业做全量刷新。查询侧看到的始终是同一张表,不感知底下用的是哪一种。

理想状态下,这个选择应该是按代价做的——同一个查询,系统自己判断增量维护和全量重算哪个更划算,甚至能证明两者等价。下面就按这个理想状态来讲,因为它才是这条路的终点。

flowchart TB DEF["<b>一份定义</b><br/>查询语义 &nbsp;+&nbsp; FRESHNESS"] DEF --> FW{"框架按 freshness<br/>选择执行形态"} FW -->|"秒级 / 分钟级"| ST["持续流式作业<br/>增量刷新"] FW -->|"小时级 / 天级"| BT["周期性批作业<br/>全量刷新"] ST --> T[("同一张表<br/>查询侧无感知")] BT --> T CHG["需求变了:<br/>只改 FRESHNESS 的值<br/>查询语义一个字不动"] -.-> DEF classDef def fill:#e3f2fd,stroke:#1976d2,color:#000 classDef exec fill:#fff3e0,stroke:#f57c00,color:#000 classDef chg fill:#f3e5f5,stroke:#7b1fa2,color:#000 class DEF def class ST,BT exec class CHG chg

对照第二部分就能看出差别:那里"这张表放 Fluss 还是只落 Paimon""补数多久跑一次""怎么和流式作业协调"是三处独立的实现决策,散落在 DDL、调度和作业配置里;这里它们收敛成定义中的一个值。

这件事最大的价值不在"改一个数字",而在它把执行细节从开发者的输入里拿掉了

四个推论值得单独说:

① 实时开发和离线开发被统一了。 你写的是一段普通 SQL,声明它要多新——不需要知道它会被跑成流式还是批式。第一部分那七个问题变成了框架的实现细节,不再是开发者的输入。这一步真正的意义是:一个只会写离线 SQL 的人,可以直接产出一张秒级的表。

② 时效性成了可修改的属性。 秒级、分钟级、T+1 不再是三种架构,而是同一份定义的三个取值。需求变了改一个值,而且语义保证不变,因为查询定义压根没动。注意因果关系:正因为执行形态被框架接管了,新鲜度才可能变成一个参数——可声明是结果,不是原因。

③ 全量刷新和增量刷新用同一份定义。 第二部分里"重建"和"流式写入"是两条独立实现的路径,靠 sequence 在存储层汇合;这里它们是同一份定义的两种执行形态,由框架决定何时用哪一种、怎么切换、怎么保证原子。那些调度依赖——补数多久跑一次、打 tag 排在补数之后——是框架内部的事了。

④ 表和作业解耦。 你操作的对象是"表",不是"那个写表的 Flink 作业"。暂停、恢复、回填、改 freshness,都是对表的操作。运维心智从"我有 37 个 Flink 作业"变成"我有 37 张表,各自有自己的新鲜度要求"。

当前能力对照。 Flink 的 Materialized Table(FLIP-435,1.20 起)已经实现了上面的骨架:一份查询定义加一个 FRESHNESS,框架自动派生执行形态和刷新流水线。但选择依据还不是代价,而是一条阈值规则——把 freshness 和 materialized-table.refresh-mode.freshness-threshold 比较,低于阈值走 CONTINUOUS(流作业),高于走 FULL(调度器周期触发批作业),也可以用 REFRESH_MODE 显式覆盖。频率同样由这个值直接映射:CONTINUOUS 下它变成作业的 checkpoint interval,FULL 下变成 cron 周期。

所以抽象是有泄漏的——名义上你声明的是"要多新",实际上你在设一个执行参数;社区在讨论让 FRESHNESS 变成可选时,理由之一就是用户得先理解这个值会变成 checkpoint interval。另外 freshness 是一个目标而非保证。功能目前仍是 MVP 阶段,作用域限于 SQL Gateway。

三、落到我们的例子

订单宽表定义一次。客服工作台要秒级 → 声明秒级 freshness,框架跑成持续流式,底下用 Fluss;运营大屏如果分钟级够用,同一份定义换个 freshness 即可。

更重要的是谁来写这张表:第一部分里它需要一个懂 watermark、懂 state、懂 retract 的人;到这一步,写它的人只需要懂业务口径和 SQL。

那些因为"实时算不起"而被推给离线的字段也一样——它们当初不进实时链路是成本问题,不是语义问题。成本消失了,"这个字段归哪条链路"这个问题本身就不存在了。

第一部分那句"时效性本该是一个参数",到这里才真正兑现。

四、再往前一步:什么才算"增量"

这里需要区分两个容易混为一谈的概念,否则展望会讲飘。

Materialized Table 统一的是入口,不是执行。 你写一份定义、声明一个 freshness,框架替你选执行形态——但底下仍然是流式作业和批作业两套机制,只是选择权从你手里交给了框架。这已经很有价值,但它没有回答一个更根本的问题:流式那份实现,是怎么保证和批式那份等价的? 答案仍然是"Flink SQL 的流批语义对齐做得足够好",而不是"从同一份定义机械推导出来的"。

真正的增量计算走的是另一条路。 以 DBSP / Feldera 为代表的一派,把增量维护形式化了:任意一个关系代数查询,都可以机械地推导出它的增量版本,且推导过程有数学保证。不存在"流式实现"和"批式实现"两份代码,只有一份定义和一个自动生成的增量程序。

这个区别对我们的意义在于——第一部分那句"逐行 diff 只能告诉你差了多少、告诉不了你该信谁",只有在后一条路上才被彻底解决。 Materialized Table 让定义只有一份,消除了实现差异;而增量推导让执行也只有一份,连"两种执行形态之间是否等价"这个问题都不再需要问。

这条路目前还远不是生产主流,工程成熟度、生态、算子覆盖面都有差距。但它指出了终点在哪:流和批不是两种计算,是同一个计算在不同触发频率下的两个特例。

五、它不解决什么

有两点必须说清楚,否则会讲成银弹。

它消除的是"不可解释",不是"不一致"。 流式结果和批式重算之间的差异依然存在——迟到数据、维表 as-of、重算边界,这些是数据本身的性质,不是实现的瑕疵。区别在于,当定义只有一份时,差异只可能来自数据到达时点,不可能来自逻辑实现不同。

它不改变物理下限。 Paimon 的分钟级下限还在,Fluss 的日志层还是必需的。FRESHNESS 是一个声明,框架据此选择执行形态,但它选得出什么形态,取决于底下有什么能力。第二部分那个地基不打好,第三部分就是空的。

顺带一句:FRESHNESS 也不会自动变便宜。声明秒级就是要付秒级的成本,只是这次成本和收益的对应关系变得显式了——你能看到"把这张表从 5 分钟提到 10 秒"意味着什么,而不是等作业跑起来才发现集群不够用。

结语

解决的问题

方式

现状

——

两条链路各自正确,靠人来协调

流批一体

重复实现、写入者竞争

一条链路,协调下沉进存储

增量语义

实时开发门槛、时效性无法声明

一份定义,执行形态交给框架

三个阶段其实在回答同一个问题的三个层次:同一份业务语义,如何在不同时效要求下保持一致。 现状的答案是"各写一遍,人来对齐";第二阶段是"写一遍,存储来对齐";第三阶段是"写一遍,而且写的人不需要知道它会被怎么执行"。

用开发成本这条线再看一遍:现状是两套逻辑两份人力,实时那套还被迫简化;第二阶段是一套逻辑,但仍要实时专家;第三阶段是一套逻辑,任何会写 SQL 的人都能产出秒级表

而每一步真正的度量,不是加了什么组件,是消失了什么:第一步消失的是两套代码和写入者竞争,第二步消失的是补数调度、tag 时机这些编排。所以这篇文章想说的其实不是某个组件好用,而是一个判断标准:

一个架构的复杂度,不在于它有几个组件,而在于有多少约束只存在于人的记忆里。

"重放的起点必须早于快照读取点"、"双流 join 成本太高,这里得改成维表 join"、"这个聚合的 state TTL 不能小于业务的最大跨度"——这类规矩一条都没写进系统里,全靠人记着。每被系统接管一条,架构就真正简单一点;反过来,加再多组件,只要它们还留在文档和某个人的脑子里,复杂度就没有下降。


我理解的流批一体:从三条链路到一个参数
https://syntomic.cn/archives/my-take-on-stream-batch-unification
作者
syntomic
发布于
2026年08月23日
更新于
2026年08月29日
许可协议