我理解的流批一体:从三条链路到一个参数
这篇文章由和claude反复对话生成。全文用同一张订单宽表贯穿,从现状的两条链路讲到 Fluss + Paimon 的改造,最后落到 Materialized Table 的增量语义。
先给出立场
"流批一体"这个词至少被用在三个不同层次上,讨论之前得先说清楚在讲哪一层:
前两层是手段,第三层才是目的。而语义之所以统一不了,根子在开发成本:实时开发的单位成本远高于离线,高到我们不得不为它单独设计一套简化的逻辑——同一张表于是长出了两条链路、两套代码、两种口径。
后面三个部分,就是把这笔成本逐层剥开:
第一部分盘现状:两条链路各自都有充分理由,但它们写同一张表,于是有了重复实现、口径分叉和写入冲突;而每一种缓解冲突的办法,都是再叠一层人工约束。
第二部分用 Fluss + Paimon 把两条链路合成一条:实时开发的成本降下来,实时链路不必再为了活下去而简化,口径分叉的根源随之消失;写入者收敛回一个,冲突在物理上不再成立。
第三部分用 Materialized Table 统一语义:一份定义加一个新鲜度声明,执行形态交给框架决定——写这张表的人不需要知道它会被跑成流式还是批式。
第一部分:现状 —— 两条链路是怎么长出来的
讲流批一体最容易犯的错误,是一上来就摆一张画满箭头的架构图,然后说"你看,Lambda 架构就是有问题"。这个讲法有个隐含前提:现在的架构是设计失误。
但真实情况是,两条链路不是一次设计出来的,而是被两个各自独立、各自成立的诉求逼出来的。先讲清楚每条链路当初为什么必须存在,后面的改造才有意义。
一、一个例子
StarRocks 上的订单宽表 dws_order_wide,主键 order_id,一张当前状态表——每个订单只保留最新状态,不按天分区。它服务客服工作台、运营大屏和即席查询。
字段就是订单的那些常规字段:order_status、pay_amt、create_time、cate_id 等等。
离线链路的加工结果落在 Hive 上,再同步到 StarRocks 对外服务;实时链路则从 Kafka 消费事件后直接写 StarRocks。两条链路是两套独立的代码,也是这张表的两个写入者。
需要回看历史状态时(同比、回溯、对账),只能另外每天存一份全量快照分区,存储按天线性增长。
二、两条链路各自的合理性
实时链路:变更必须马上可见。 订单状态变更要秒级出现在客服工作台。用户打电话说"我付款了状态怎么还没变",客服看到的必须是当前态。天级满足不了,小时级也满足不了。这是业务硬约束,不是选型偏好。
但它的逻辑是刻意简化过的:只读事件流,再点查 HBase 把维度打宽,不做多流 join。核心矛盾在双流 join——两边都要存全量 state,状态规模压不住;而且任何一边的更新都会触发回撤,而回撤是"先撤旧值、再发新值"两条消息,这中间的一瞬间下游看到的就是错的:行短暂消失,或者聚合值先掉下去再涨回来。查询随时都在发生,正好落在这个窗口上就会读到错值,事后还查不到、复现不了。为了让这条链路能被维护,我们在逻辑上做了一些简化。
离线链路:需要一份会收敛的全量。 实时链路的正确性依赖太多外部条件——Kafka 可能丢消息、作业可能挂掉丢状态、CDC 可能漏采、上线可能引入 bug。而这些问题在实时链路里没有自愈机制:一行写错了就一直错下去。所以必须有一条从上游全量快照出发、可重跑、可校验的链路,每天把整表覆盖一遍。它的价值不是时效,是收敛性:无论昨天发生了什么,第二天早上数据一定是对的。这条链路用 Spark SQL 做分层加工,跑的是天级批任务。
两个诉求:低延迟、正确性。彼此独立,各自成立。
(现实里往往还有第三条小时级链路,用来刷新那些"T+1 太慢、但上游只能到小时"的指标。它本质上是离线链路换了个调度周期,下面所有的论证对它同样适用,为了聚焦这里就不单列了。)
三、然后问题出现了
两条链路各自合理,写同一张表同样合理——使用方要的就是一张能直接查的宽表,不该也没法感知某个字段从哪条链路来。所以问题不在"为什么写同一张表",而在于:两个写入者、两种写入语义、两种时效,落在同一张表上,却没有任何机制协调它们。
1. 两套逻辑、六个组件,开发成本高
最直接的成本是同一份业务语义写了两遍:Flink SQL 一遍、Spark SQL 一遍,两种引擎方言。改一次口径要改两处,漏改一处就是长期存在、难以发现的不一致。
但更麻烦的是这两遍并不等价。离线链路能做复杂的多表 join,实时链路只能点查打宽——所以有些字段在实时侧算起来代价太高,索性交给离线补;有些字段两边都算得出,但走的是完全不同的路径。两边的差异不是 bug,是设计使然。 也正因如此,两条链路的结果永远对不齐——你没法要求两套本来就不同的逻辑给出一样的结果。
组件账也不轻:Kafka、HBase、Flink、Spark、Hive、StarRocks——一个订单宽表的字段,要经过六个系统。每加一个字段,都要判断它归哪条链路、会不会和另一条冲突、维表要不要同步扩一列。
而两条链路里,实时那条的单位成本远高于离线:
离线开发是写一段 SQL;实时开发是写一段 SQL,再回答七个关于执行的问题。 所以同一个人做实时的产出可能只有做离线的三分之一,而能做实时的人本来就少。
这就是上面那个"刻意简化"的由来——不是能力选择,是成本倒逼。而简化的代价,就是口径先天分叉。
2. 全量覆盖的丢失窗口
上面是成本问题,下面是正确性问题。
离线任务凌晨 01:00 读取上游全量快照,此刻订单 O1 是"未支付";跑批、加工、同步走了两小时;02:00 用户付款,实时链路已经把 StarRocks 写成"已支付";03:00 离线结果整表覆盖,把它退回"未支付"。
覆盖是整表的,所以影响范围不是某一类订单,而是读取快照到覆盖完成这段窗口内所有发生过变更的行。这些行会停在旧值上,直到它下一次发生业务变更才被实时链路带回正确状态——如果它已经进入终态、之后长期不变,那就是长期错误。
凌晨订单量低,绝对数量看着不大,而且这些行"迟早会自己好",所以它长期不被当成一个问题——直到某天客服反馈某个订单状态一直不对。
根因不是覆盖本身,而是胜负由到达顺序决定,而不是由数据版本决定。这是典型的 last-write-wins by arrival time:只要存在多个写入者、且各自数据版本时刻不同,就必然丢数据。
3. 每个解法,都是再加一层人工维护的约束
丢失窗口有两条解法:加版本列,或者重放事件。
解法一:版本列。 给结果表指定一个 sequence 字段——取业务行自身的 update_time——让引擎按版本大小定胜负,而不是按到达顺序。离线覆盖携带的是快照时刻的旧版本,天然写不进那些已经被流式更新过的行。
机制很简单,但它有三个前提,每一个都不由数仓控制:
一是存储得支持。 StarRocks 的主键表可以指定 sequence 列,MySQL 这类只认"后写覆盖先写"的存储就没有这个能力——换个 sink,这条路直接走不通。
二是版本字段没得选。 处理时间等于没做,又回到按到达顺序;CDC 的位点或采集时间离线侧拿不到,两个写入者之间不可比。唯一两边都能拿到、且语义同源的就是业务行自身的 update_time。
三是这个字段的可靠性由上游决定。 触发器直接改表、DBA 手工修数据、ORM 的选择性更新——任何一条不更新 update_time 的写入路径,都会让版本比较失效。而且是静默失效:不报错,只是某些行的值一直不对。
还有一个次要代价:写入从覆盖退化成 merge,行集合不再严格对齐。上游物理删除、而流式恰好漏掉 -D 的那些行会永远留着——覆盖时它们本来会被顺手清掉。
解法二:重放事件。 不加版本列,而是在离线覆盖完成后触发实时作业重启,把 offset 拨回快照读取点之前,让窗口内的事件重新走一遍。因为写入者始终只有一个、事件按顺序依次重放,最终值必然正确。
问题出在"最终"这两个字上——catch-up 期间表处于一个对外可见的错误中间态:它要按顺序把几小时前的事件重新放一遍,这段时间查询会看到数据先退回旧值、再逐步恢复,而且不同行的恢复进度还不一样。所以重放不能直接在对外的那张表上做,需要配一套双链路切换:覆盖和重放都在备表上完成,追平并校验通过后再把读流量切过去。
代价是双份存储、一层切换编排,以及对作业的硬要求:必须无状态,或状态可从重放中重建——一旦有跨天累计、全局去重、双流 join,重置 offset 就是丢 state,结果直接算错。
两条路都有效,很多团队就这么跑了好几年。但把它们放在一起看,是同一个模式:每一次缓解,都是在架构之上再叠一层需要人工维护的约束。 update_time 必须被上游可靠维护、重放起点必须早于快照读取点、切换前后两张表的保留期要对齐、有状态的作业不能用重放——这些约束没有一条写进任何系统里,只存在于文档、注释和某个人的记忆中。正确性依赖的是"你记得所有约束",而不是"系统保证了它们"。
4. 时效性被硬编码进了架构
现在这条边界是显式可见的:秒级的走 Kafka,天级的走 Spark。 哪个字段要多快,决定的不是一个参数,而是它归哪条链路、用哪套代码、走哪套调度。
如果哪天某个指标要从天级提到分钟级,或者反过来实时链路成本太高要降到五分钟——答案都是把字段在链路之间搬家,重写、重配、重新验证。
时效性本该是一个参数,现在却是一个架构决策。 这不是某个组件的能力缺陷,是整套架构的表达能力缺陷。
四、小结
两条链路各有充分理由,已有的工程实践也确实压住了大部分问题。现状不是不能用。
但两件根本的事从头到尾没被解决:
开发成本:同一份语义写两遍、跨六个组件,而且实时那遍因为成本太高被迫简化,导致口径先天分叉、差异不可解释。
协调成本:两个写入者落在同一张表上,冲突全靠人来编排,而这些编排只存在于文档和某个人的记忆里。
这两条其实是同一件事的两面:如果实时链路不必再简化,它就不需要离线那条来补齐字段;如果只剩一个写入者,冲突也就无从谈起。 所以下一步的着力点是让实时开发不再那么贵——链路合成一条是这件事的结果,不是先定下的目标。
第二部分:Fluss + Paimon —— 把协调下沉进存储
两条链路各有各的诉求,硬合只是把问题挪个地方。所以先看清每个诉求该由谁承接。
一、职责重新划分
关键在第一行:秒级不再需要一条独立链路,它变成同一张表的一个存储层。 这是整个阶段二的支点,后面所有收益都从这一条推出来。
二、为什么需要 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 那半边完全成立,上面四条能力换不来一个新组件的运维成本。我们引入它,是因为客服工作台的秒级诉求是硬的,而它一旦引入,实时侧就顺势建成了真正的分层数仓,而不只是把一条流写进湖。
三、架构
四个变化点:
分层第一次同时存在于流和批。 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 不安全:
原因在于 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 在离线快照读取之后、覆盖完成之前发生了支付。
没有任何一个时刻会退回旧值。
第一部分那张丢失窗口的时序图,在这里没有对应物。这是整个改造最值得记住的一句话:问题不是被修好了,是在物理上不成立了。
六、逐条对账第一部分的问题
前几条都出自同一个原因:协调机制从人手里下沉到了存储层。 现状里那些约束——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 的人来做,团队产能没有变化。这是整笔成本账里最贵的一块,也是这一节要谈的东西。
这一节是展望,不是方案——相关能力还在演进中,这里讲的是方向。
一、第二部分留下了什么
除了实时开发这条,其余的约束也只是换了形态:
这六条比第一部分那些好得多——它们至少是显式配置,可以被 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 和一个"要多新"。哪天大屏说分钟级够用了,改的就是那个 INTERVAL;REFRESH_MODE 也可以显式写成 CONTINUOUS 或 FULL 来覆盖框架的判断。
框架根据 freshness 决定执行形态:要求秒级或分钟级,跑成一个持续的流式作业做增量刷新;要求小时或天级,退化成周期性的批作业做全量刷新。查询侧看到的始终是同一张表,不感知底下用的是哪一种。
理想状态下,这个选择应该是按代价做的——同一个查询,系统自己判断增量维护和全量重算哪个更划算,甚至能证明两者等价。下面就按这个理想状态来讲,因为它才是这条路的终点。
对照第二部分就能看出差别:那里"这张表放 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 不能小于业务的最大跨度"——这类规矩一条都没写进系统里,全靠人记着。每被系统接管一条,架构就真正简单一点;反过来,加再多组件,只要它们还留在文档和某个人的脑子里,复杂度就没有下降。