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

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

先给出立场

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

层次

主张

代表

引擎层

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

Flink 流批统一的 API 与运行时

存储层

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

Paimon / 湖仓一体

语义层

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

本文的主题

前两层是手段,第三层才是目的。引擎统一了、存储统一了,如果口径还是各写一遍、时效还是硬编码在架构里,流批一体就没有真正发生。

所以我理解的流批一体是一句话:

同一份业务语义,在不同时效要求下保持一致;而"要多新"应该是一个可声明的参数,不是一个架构决策。

这句话同时定义了目标和检验标准。下面三个部分,都是在逼近它。

状态

对那句话的回答

现状

三条链路,三套代码

各写一遍,人来对齐

流批一体改造

一条链路,Fluss + Paimon

写一遍,存储来对齐

增量语义

Materialized Table

写一遍,时效是它的一个参数

三个阶段不是三种技术的堆叠,是同一个问题被逐层剥开:先消灭重复实现,再消灭执行形态的选择。

而衡量每一步进展的单位,也不是"引入了什么组件",而是:少了多少只存在于人记忆里的约束。

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

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

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

一、一个例子

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

字段按来源分两组:

字段

来源

时效

A · 事实主干

order_statuspay_amtcreate_timecate_id

实时链路

秒级

C · 时效敏感指标

user_levelrisk_score

小时链路

小时级

底层存储部分上了湖:ODS 和离线的加工结果落在 Paimon 上,再同步到 StarRocks 对外服务;实时链路和小时链路则各自直接写 StarRocks。三条链路是三套独立的代码,也是这张表的三个写入者。

需要回看历史状态时(同比、回溯、对账),只能另外每天存一份全量快照分区——存储按天线性增长,或者做拉链表,逻辑复杂且边界条件极易出错。

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

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

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

小时链路:离线是天级,某些指标等不了。 用户等级、风控分这类字段离线也会产出,但 T+1 太慢——风控分变了当天就得生效。可它们来自算法侧,模型按小时批量跑,也做不到秒级。需求要求更快,上游能力封顶在小时,于是在实时和离线之间插了第三条链路:批读算法侧的产出,用部分列更新直接写 StarRocks。它不进湖——数据只是几个字段的刷新,为它单独走一遍湖上的分层加工没有意义。用部分列是因为这条链路手里只有 C 组的值,整行写入会把实时链路刚写进去的 A 组状态覆盖掉。

三个诉求:低延迟、正确性、时效错配。彼此独立,各自成立。

flowchart LR DB[("MySQL")] -->|CDC| MQ[("Kafka")] MQ --> F1["Flink SQL<br/>实时链路"] MQ --> ODS["Paimon ODS"] ALG[("算法侧<br/>小时产出")] --> HR["Flink SQL<br/>小时链路"] ODS -->|Spark SQL 天级| DWS["Paimon DWD / DWS"] F1 -->|"A 组 · 整行 upsert · 秒级"| SR HR -->|"C 组 · 部分列更新 · 小时级"| 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 classDef hr fill:#f3e5f5,stroke:#7b1fa2,color:#000 class F1 rt class ODS,DWS off class HR hr

三、为什么秒级不能也建在湖上

数据都已经在 Paimon 上了,为什么不把实时链路也建在湖上,用 Flink 流式写 Paimon,省掉那条独立的 Kafka 链路?

因为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 分钟。这是"周期性提交不可变快照"这个模型的固有下限。

秒级场景只能绕开湖,单独拉一条 Kafka + Flink SQL 直写 StarRocks 的链路。这条链路的存在理由只有这一条。

四、然后问题出现了

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

1. 离线被迫成为另外两条的超集

覆盖是整表粒度的。如果离线只产出 A 组字段,小时链路写进去的 C 组字段每天都会被冲成 NULL。为了避免这一点,离线必须也把 C 组一起产出。而 A 组的计算逻辑,本来就和实时链路完全重复。

结果是:离线链路重新实现了另外两条的全部逻辑。代码重复率不是 1/3,是接近 100%,而且还是 Spark SQL 和 Flink SQL 两种方言。改一次口径要改三处,漏改一处就是长期存在、难以发现的不一致。

2. C 组字段会被天级快照退回

C 组的需求很明确——要的就是最新状态,风控分变了就该立刻反映出来。小时链路每小时刷一次,符合预期。但离线链路产出的是天级快照,覆盖时会把这个字段退回到前一天的值,直到下一次小时任务跑完才恢复。

这中间是一段确定会出现、但没人盯着的脏数据窗口。它比下面要讲的丢失窗口温和——小时链路会自愈——但根因完全一样:携带旧版本的写入者后到,于是旧值赢了。

3. 全量覆盖的丢失窗口

同样的机制,换个字段就可能是长期性的。

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

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

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

凌晨订单量低,绝对数量看着不大,而且这些行"迟早会自己好",所以它长期不被当成一个问题。但它绕过了所有常规对账规则:金额总量偏差远低于任何合理阈值,逐行 diff 又要在覆盖之后才能做——那时两边已经一致了,因为实时链路刚被覆盖成了离线的值。

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

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

丢失窗口有两条解法。

加版本列,让引擎按 sequence 大小定胜负而不是按到达顺序。这在全量表场景下其实相当契合——离线快照携带的是旧版本,天然写不进已有更新的行。代价是写入退化成 merge,失去了整表覆盖的删除兜底能力:上游已物理删除、而实时链路漏掉删除事件的行,永远清不掉,需要额外一条低频的 anti-join 清理任务。而且它静默依赖 update_time 的可靠性——一旦上游存在不更新该字段的变更路径,比较会失效且无告警。

重放事件,离线覆盖完成后触发实时作业重启、offset 回到读取点之前。因为写入者只有一个,事件按序重放,最终值必然正确,不需要版本列,也保住了删除兜底。但它有一个硬门槛:作业必须无状态,或状态可从重放中重建——一旦有跨天累计、全局去重、双流 join,重置 offset 就是丢 state,结果直接算错。

再加上"离线接管整表"这个动作本身也要保证原子、可校验、可回滚,于是还要一层影子表加对账门禁:离线结果先写进一张对外不可见的影子表,此时新旧两份数据同时存在,跑一轮对账——行数偏差、主键是否重复、金额总量偏差、按主键逐行 diff 各字段——通过了才原子换表,没通过就不切、告警、维持实时结果。

这套机制本身是对的,顺序也很关键:门禁必须在换表之前,先覆盖再校验就已经没有退路了。但它引入的东西不少——一份冗余的影子表存储、一套随字段数线性增长的对账规则、以及三处必须彼此对齐的保留期(影子表保留、Kafka retention、上游快照保留)。这三个保留期一旦错配,回滚就会失败,而没有任何系统会帮你检查这件事。

更关键的是它解决不了最根本的那个问题:逐行 diff 只能告诉你"差了 3%“,告诉不了你"该信谁”。因为两边是两套代码、两种引擎,差异既可能是 bug,也可能是本来就不同的口径。所以阈值只能设得很宽,最终退化成一个"防止灾难"的兜底,而不是"保证正确"的证明。

这些做法都有效,很多团队就这么跑了好几年。但把它们放在一起看,是同一个模式:每一次缓解,都是在架构之上再叠一层需要人工维护的约束。 重放起点要早于读取点、三处保留期要对齐、对账阈值要随字段调整、状态型作业不能用重放——这些约束没有一条写进任何系统里,只存在于文档、注释和某个人的记忆中。正确性依赖的是"你记得所有约束",而不是"系统保证了它们"。

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

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

如果哪天算法侧说风控分能做到分钟级了,或者反过来实时链路成本太高要降到五分钟——答案都是把字段在链路之间搬家,重写、重配、重新对账。

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

五、小结

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

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

  • 同一份业务语义在三条链路里各写了一遍,差异不可解释。

  • 时效性是架构的一部分,而不是可声明的参数。

接下来两个阶段正是沿着这两条线走。Fluss 补上 Paimon 缺的那个日志层——写入即可见,不依赖 checkpoint 做可见性,端到端百毫秒级;而它和 Paimon 是同一张表的热层与冷层,由 tiering 按 offset 顺序单向归档、查询做 union read。同样是承接秒级,Kafka 意味着另起一条链路、另一套 SQL、另一份口径,Fluss 意味着同一张表多了一个时效档位,分层不变、SQL 不变、口径不变

Materialized Table 补齐最后一层语义:一份定义,FRESHNESS 决定它退化成哪种执行形态——秒级、分钟级还是 T+1,不再是链路选择,而是一个参数。

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

第一部分的结论是两条:同一份语义在三条链路里各写了一遍;时效性是架构决策而不是参数。

这一部分解决第一条。方法不是"把三条链路合并"——那只是把问题挪个地方——而是先看清每个诉求该由谁承接

一、职责重新划分

诉求

现状由谁承担

改造后由谁承担

秒级可见

独立的 Kafka 链路

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

收敛性兜底

Spark 天级整表覆盖

全量快照重灌 + sequence 版本比较

时效错配的中间档

独立作业直接写 StarRocks

同一张表的 partial-update merge engine

历史回溯

每天另存全量快照分区 / 拉链表

Paimon tag 时间旅行

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

二、为什么是 Fluss,而不是继续用 Kafka

这个问题必须先回答,否则整套方案会显得是在堆技术名词。

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

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

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

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

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

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

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

三、架构

flowchart LR DB[("MySQL")] -->|CDC 增量| ODS["Fluss ODS"] SNAP[("上游全量快照")] -->|每日重灌| ODS ODS --> DWD["Fluss DWD"] --> DWS["Fluss DWS<br/>主键表 · sequence"] ALG[("算法侧<br/>小时产出")] -->|"C 组 · partial update"| DWS 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 表的另一组列,由 partial-update merge engine 合并——从"三个写入者用不同语义互相避让"变成"一张表声明了它如何合并多路输入"。

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. partial-update 怎么合并多路输入

宽表声明 merge engine 为 partial-update,并指定 sequence 字段。A 组的流只写 A 组的列,C 组的小时作业只写 C 组的列,存储层按主键合并、按 sequence 决定同列的新旧。

它直接解释了"C 组被天级快照退回"为什么消失:日终快照和小时快照不再是两个覆盖动作在竞争,而是同一张表上按 sequence 排序的两个版本,旧的那个不可能赢。现状里"离线不要覆盖 C 组"只能靠离线作业自觉把 C 组一起产出来实现,改造后这件事由合并规则保证。

约束没有消失:各路输入的字段组不能重叠,sequence 必须可比。

3. 兜底:从"覆盖"变成"重灌"

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

因为目标是一张当前状态表,这件事被大幅简化了——不需要决定"重算哪几天",也没有重算窗口,只需要定期把上游全量快照重新灌进同一张表,每行携带自己的 update_time 作为 sequence。

于是收敛靠版本比较完成:已经有更新版本的行,旧快照写不进去;流式漏掉或写错的行,被快照补正。胜负由数据版本决定,不由到达顺序决定。

这和现状是两回事。现状里根本没有版本列——三个写入者靠"各写各的列"这个人为约定来避让:实时写 A 组,小时链路部分列写 C 组,离线整表覆盖全部。约定不在任何地方被声明,而整表覆盖恰好打破了它,于是就有了第一部分那些问题。sequence 在第一部分只是作为一个候选解法被提出来,并没有落地。

改造后它是 Fluss 表定义的一部分。表声明了"我如何合并输入",这个语义随表走——union read、tiering 归档、批读、tag 时间旅行,看到的都是同一套合并规则;而且它作用在数据源头,整条链路都被保护,不是只在终点加一道防线。

更本质的变化是冲突消解的性质变了

整表覆盖

带 sequence 的重灌

语义

“我这一份就是全部真相”

“我这里有若干行的某个版本”

能否和别人的写入合并

不能

正确性靠什么

单一写入者,或人工编排时序

合并规则本身

覆盖是一个不可合并的操作——它不携带版本信息,也不承认还有别的真相来源,所以两个写入撞在一起时系统没有依据判断谁对,只能按谁后到。这就是为什么现状必须靠"消灭第二个写入者"或"重放 offset"这类外部编排来兜底。

重灌是可合并的:按 sequence 取最大这个规则幂等、可交换、可结合。同一行写十遍结果一样,谁先到谁后到结果一样。最终状态由数据本身唯一确定,与写入者数量和到达顺序都无关。

所以第一部分那句"只要存在多个写入者就必然丢数据",隐含前提是"冲突靠到达顺序消解"。前提一旦去掉,结论也就不成立了——多写入者从来不是问题的根源,"用到达顺序决定胜负"才是。

只用 update_time 做 sequence 是不够的

这里有个容易漏掉的坑:sequence 被要求同时干两件事——挡住陈旧快照,以及让修正生效。而单靠业务 update_time,第二件做不到。

原因是重灌读的是同一份上游数据,update_time 不变。所以当流式把某行算错时,重灌携带的版本和表里那行是相等的,不是更大的:

场景

行是否已存在

update_time 比较

只用 update_time

流式漏采整行

不存在

——

能修(插入即生效)

流式把值算错了

存在

相等

修不了

快照比流式旧

存在

更小

不生效(正确行为)

中间那行恰恰是兜底最需要覆盖的场景。解法是复合 sequence:把重灌批次号作为次级比较字段,流式写入固定为 0,每次重灌用单调递增的批次号。

比较规则变成先比 update_time、相等再比 reload_batch,于是三种场景各归其位:

sequenceDiagram participant RT as 流式链路 participant T as Fluss 表 participant RL as 重灌作业 Note over T: 源行版本 update_time = 10:00 RT->>T: 写入 (10:00, batch=0) Note over T: 值算错了 ✗ rect rgb(240, 250, 240) RL->>T: 重灌 (10:00, batch=7) Note over T: update_time 相等<br/>batch 7 > 0 → 重灌赢<br/>值被修正 ✓ end Note over T: 源行发生真实变更 → 11:00 rect rgb(232, 244, 253) RT->>T: 写入 (11:00, batch=0) Note over T: update_time 更大<br/>→ 流式赢,秒级生效 ✓ end rect rgb(255, 248, 240) RL->>T: 次日重灌,快照仍是 (10:00, batch=8) Note over T: update_time 更小<br/>→ 直接丢弃 ✓ end

这样把"同一个源版本上,重灌是权威值、流式是它的低延迟近似"这条规则显式写进了表定义,而不是依赖引擎在 sequence 相等时的默认 tie-break 行为。partial-update 场景下每个列组各有自己的 sequence group,分别按这个规则比较,互不干扰。

逻辑错了也用同一个机制。 口径变更、新增字段、算法写错——改完 SQL 之后做一次全量重灌就行:目标是一张当前状态表,每行的正确值完全由上游那行的当前值决定,重灌会让所有行用新逻辑重算一遍,靠更大的 reload_batch 覆盖旧值。不需要另起一张新表,也不需要 catalog 切换。

反过来则不成立:重放日志替代不了重灌。CDC 漏采的消息从来没进过日志,重放也变不出来,只能靠上游快照补。所以这套架构里只需要一个修正机制,就是重灌。

例外只有一种:如果某个字段的计算依赖变更历史而不只是当前值(累计次数、首单标记、状态机变迁轨迹),重灌算不出来,那部分才需要从日志重放重建。我们这张宽表没有这类字段。

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

现状要保留历史,无非两条路——每天存一份全量分区(存储按天线性增长),或者做拉链表(逻辑复杂、边界条件多、极易出错)。

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

这里有一个必须讲清楚的区分——在 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 两种执行模式

C 组被天级快照退回

解决

partial-update + sequence,表级语义

全量覆盖的丢失窗口

消失

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

层层人工约束

大幅缓解

影子表和对账门禁不再需要;水位收敛成 offset 和 tag

历史回溯要额外存全量

解决

tag 增量存储

时效性硬编码

没解决

见第七节

前几条都出自同一个原因:协调机制从人手里下沉到了存储层。 现状里那些约束——重放起点要早于读取点、三处保留期要对齐、对账阈值要随字段调整、状态型作业不能重放——要么消失,要么变成表属性和调度依赖这种可以被系统检查的东西。

七、必须诚实讲的五件事

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

删除兜底需要单独一条任务。 全量重灌补不了"上游已物理删除、而流式漏掉了删除事件"这种情况——快照里没有这行,重灌就是个 no-op。需要一条低频(周级)的 anti-join:比对快照主键集合和表内主键集合,对差集发出删除事件。

sequence 的可靠性是整个兜底的地基。 首先它必须在参与竞争的写入者之间可比——都取业务行自身的 update_time,而不是各自的处理时间。其次,如果上游存在不更新 update_time 的变更路径(触发器、DBA 直接改库、部分 ORM 的选择性更新),版本比较会静默失效。最后,复合 sequence 里的重灌批次号必须严格单调递增,这条要由调度保证。这条在当前状态表场景下尤其要紧,因为兜底完全建立在它之上——上线前必须实测验证,不能假设。

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

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

关于替换 Kafka:本文只考虑数仓链路自己的消费方。迁移期建议按 topic 分批切,每批切完立刻下线 Kafka 侧,不要长期双跑——并存期本身就是又一次多数据源,正是我们要消灭的东西。

八、一句话结论

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

整套设计最终收敛成三件事,各自只有一个机制:秒级当前状态靠 union read,历史回溯靠 tag,数据和逻辑的修正都靠重灌加 sequence。没有冷热分离,没有视图水位,没有影子表门禁,也不需要为改口径另起一张表。

留下的是"何时刷新、要多新"这个决策——它还在人的脑子里、还在调度配置里,没有被声明出来。这就是第三部分要处理的东西。

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

第二部分把三条链路收成了一条,协调机制下沉进了存储。但还剩一个问题没动:时效性仍然是架构决策,不是参数。

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

一、第二部分留下了什么

改造之后,架构确实干净了:三件事,三个机制。但如果把"还需要人来维护的约束"列一遍,会发现它们并没有消失,只是换了形态:

约束

现在存在于哪里

谁来保证

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

建表 DDL

人的判断

重灌多久跑一次

调度配置

人的判断

重灌批次号必须单调递增

调度实现

人的实现

打 tag 必须排在重灌之后

调度依赖

人配的依赖

tag 保留多久

表属性 + 清理任务

人的判断

sequence 字段是否可靠

上游的写入习惯

没人保证

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

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

这就是最后那个缺口:需求变化时——算法侧提速了、某个大屏要求更快、某条链路成本太高要降级——你要做的还是老三样:改代码、改配置、重新对账。第一部分我们说"时效性本该是一个参数",第二部分换掉了地基,但那个参数依然不存在

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

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

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

框架根据 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、调度和作业配置里;这里它们收敛成定义中的一个值。

三个推论值得单独说:

① 时效性成了可修改的属性。 秒级、分钟级、T+1 不再是三种架构,而是同一份定义的三个取值。需求变了,改的是一个值,而不是一条链路——而且改完之后语义保证不变,因为查询定义压根没动。

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

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

三、落到我们的例子

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

C 组字段那件事最能说明问题:现状里它因为"上游只能小时级"而被迫成为独立链路,第二部分里它变成同一张表的一个列组,但仍然要单独配作业和调度。到这一步,它就是声明里的一个 freshness 值——算法侧哪天提速到分钟级,改这个值即可,不需要把字段在链路之间搬家,也不需要重新对账。

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

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

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

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

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

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

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

五、它不解决什么

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

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

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

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

六、三个阶段的一条线

解决的问题

方式

现状

——

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

流批一体

重复实现、多写入者竞争

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

增量语义

时效性无法声明

一份定义,freshness 是参数

三个阶段其实在回答同一个问题的三个层次:同一份业务语义,如何在不同时效要求下保持一致。

现状的答案是"各写一遍,人来对齐";第二阶段是"写一遍,存储来对齐";第三阶段是"写一遍,时效是它的一个参数"。

值得注意的是每一步消失的东西:第一步消失的是三套代码和多写入者竞争,第二步消失的是重灌调度、tag 时机、批次号单调这些编排。架构演进的度量单位不是加了什么组件,而是少了多少需要人记住的东西。

结语

这篇文章真正想说的其实不是某个组件好用,而是一个判断标准:

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

第一部分那些重放起点、保留期对齐、对账阈值,第二部分剩下的重灌节奏和 tag 时机——每一条被系统接管,架构就真正简单了一点。反过来,加再多组件,只要约束还在文档和某个人的脑子里,复杂度就没有下降。


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