我理解的流批一体:从三条链路到一个参数
这篇文章由和claude反复对话生成。全文用同一张订单宽表贯穿,从现状的三条链路讲到 Fluss + Paimon 的改造,最后落到 Materialized Table 的增量语义。
先给出立场
"流批一体"这个词被用在三个不同层次上,讨论之前必须先说清楚在讲哪一层:
前两层是手段,第三层才是目的。引擎统一了、存储统一了,如果口径还是各写一遍、时效还是硬编码在架构里,流批一体就没有真正发生。
所以我理解的流批一体是一句话:
同一份业务语义,在不同时效要求下保持一致;而"要多新"应该是一个可声明的参数,不是一个架构决策。
这句话同时定义了目标和检验标准。下面三个部分,都是在逼近它。
三个阶段不是三种技术的堆叠,是同一个问题被逐层剥开:先消灭重复实现,再消灭执行形态的选择。
而衡量每一步进展的单位,也不是"引入了什么组件",而是:少了多少只存在于人记忆里的约束。
第一部分:现状 —— 三条链路是怎么长出来的
讲流批一体最容易犯的错误,是一上来就摆一张画满箭头的架构图,然后说"你看,Lambda 架构就是有问题"。这个讲法有个隐含前提:现在的架构是设计失误。
但真实情况是,三条链路不是一次设计出来的,而是被三个各自独立、各自成立的诉求逼出来的。先讲清楚每条链路当初为什么必须存在,后面的改造才有意义。
一、一个例子
StarRocks 上的订单宽表 dws_order_wide,主键 order_id,一张当前状态表——每个订单只保留最新状态,不按天分区。它服务客服工作台、运营大屏和即席查询。
字段按来源分两组:
底层存储部分上了湖:ODS 和离线的加工结果落在 Paimon 上,再同步到 StarRocks 对外服务;实时链路和小时链路则各自直接写 StarRocks。三条链路是三套独立的代码,也是这张表的三个写入者。
需要回看历史状态时(同比、回溯、对账),只能另外每天存一份全量快照分区——存储按天线性增长,或者做拉链表,逻辑复杂且边界条件极易出错。
二、三条链路各自的合理性
实时链路:变更必须马上可见。 订单状态变更要秒级出现在客服工作台。用户打电话说"我付款了状态怎么还没变",客服看到的必须是当前态。天级满足不了,小时级也满足不了。这是业务硬约束,不是选型偏好。
离线链路:需要一份会收敛的全量。 实时链路的正确性依赖太多外部条件——Kafka 可能丢消息、作业可能挂掉丢状态、CDC 可能漏采、上线可能引入 bug。而这些问题在实时链路里没有自愈机制:一行写错了就一直错下去。所以必须有一条从上游全量快照出发、可重跑、可校验的链路,每天把整表覆盖一遍。它的价值不是时效,是收敛性:无论昨天发生了什么,第二天早上数据一定是对的。这条链路用 Spark SQL 做分层加工,跑的是天级批任务。
小时链路:离线是天级,某些指标等不了。 用户等级、风控分这类字段离线也会产出,但 T+1 太慢——风控分变了当天就得生效。可它们来自算法侧,模型按小时批量跑,也做不到秒级。需求要求更快,上游能力封顶在小时,于是在实时和离线之间插了第三条链路:批读算法侧的产出,用部分列更新直接写 StarRocks。它不进湖——数据只是几个字段的刷新,为它单独走一遍湖上的分层加工没有意义。用部分列是因为这条链路手里只有 C 组的值,整行写入会把实时链路刚写进去的 A 组状态覆盖掉。
三个诉求:低延迟、正确性、时效错配。彼此独立,各自成立。
三、为什么秒级不能也建在湖上
数据都已经在 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 离线结果整表覆盖,把它退回"未支付"。
覆盖是整表的,所以影响范围不是某一类订单,而是读取快照到覆盖完成这段窗口内所有发生过变更的行。这些行会停在旧值上,直到它下一次发生业务变更才被实时链路带回正确状态——如果它已经进入终态、之后长期不变,那就是长期错误。
凌晨订单量低,绝对数量看着不大,而且这些行"迟早会自己好",所以它长期不被当成一个问题。但它绕过了所有常规对账规则:金额总量偏差远低于任何合理阈值,逐行 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 —— 把协调下沉进存储
第一部分的结论是两条:同一份语义在三条链路里各写了一遍;时效性是架构决策而不是参数。
这一部分解决第一条。方法不是"把三条链路合并"——那只是把问题挪个地方——而是先看清每个诉求该由谁承接。
一、职责重新划分
关键在第一行:秒级不再需要一条独立链路,它变成同一张表的一个存储层。 这是整个阶段二的支点,后面所有收益都从这一条推出来。
二、为什么是 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 那半边完全成立,上面四条能力换不来一个新组件的运维成本。我们引入它,是因为客服工作台的秒级诉求是硬的,而它一旦引入,实时侧就顺势建成了真正的分层数仓,而不只是把一条流写进湖。
三、架构
四个变化点:
分层第一次同时存在于流和批。 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 时间旅行,看到的都是同一套合并规则;而且它作用在数据源头,整条链路都被保护,不是只在终点加一道防线。
更本质的变化是冲突消解的性质变了:
覆盖是一个不可合并的操作——它不携带版本信息,也不承认还有别的真相来源,所以两个写入撞在一起时系统没有依据判断谁对,只能按谁后到。这就是为什么现状必须靠"消灭第二个写入者"或"重放 offset"这类外部编排来兜底。
重灌是可合并的:按 sequence 取最大这个规则幂等、可交换、可结合。同一行写十遍结果一样,谁先到谁后到结果一样。最终状态由数据本身唯一确定,与写入者数量和到达顺序都无关。
所以第一部分那句"只要存在多个写入者就必然丢数据",隐含前提是"冲突靠到达顺序消解"。前提一旦去掉,结论也就不成立了——多写入者从来不是问题的根源,"用到达顺序决定胜负"才是。
只用 update_time 做 sequence 是不够的
这里有个容易漏掉的坑:sequence 被要求同时干两件事——挡住陈旧快照,以及让修正生效。而单靠业务 update_time,第二件做不到。
原因是重灌读的是同一份上游数据,update_time 不变。所以当流式把某行算错时,重灌携带的版本和表里那行是相等的,不是更大的:
中间那行恰恰是兜底最需要覆盖的场景。解法是复合 sequence:把重灌批次号作为次级比较字段,流式写入固定为 0,每次重灌用单调递增的批次号。
比较规则变成先比 update_time、相等再比 reload_batch,于是三种场景各归其位:
这样把"同一个源版本上,重灌是权威值、流式是它的低延迟近似"这条规则显式写进了表定义,而不是依赖引擎在 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 不安全:
原因在于 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 在离线快照读取之后、覆盖完成之前发生了支付。
没有任何一个时刻会退回旧值。
第一部分那张丢失窗口的时序图,在这里没有对应物。这是整个改造最值得记住的一句话:问题不是被修好了,是在物理上不成立了。
六、逐条对账第一部分的问题
前几条都出自同一个原因:协调机制从人手里下沉到了存储层。 现状里那些约束——重放起点要早于读取点、三处保留期要对齐、对账阈值要随字段调整、状态型作业不能重放——要么消失,要么变成表属性和调度依赖这种可以被系统检查的东西。
七、必须诚实讲的五件事
时效性依然是硬编码的。 你还是要决定哪张表建在 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。没有冷热分离,没有视图水位,没有影子表门禁,也不需要为改口径另起一张表。
留下的是"何时刷新、要多新"这个决策——它还在人的脑子里、还在调度配置里,没有被声明出来。这就是第三部分要处理的东西。
第三部分:增量语义 —— 一点展望
第二部分把三条链路收成了一条,协调机制下沉进了存储。但还剩一个问题没动:时效性仍然是架构决策,不是参数。
这一节是展望,不是方案——相关能力还在演进中,这里讲的是方向。
一、第二部分留下了什么
改造之后,架构确实干净了:三件事,三个机制。但如果把"还需要人来维护的约束"列一遍,会发现它们并没有消失,只是换了形态:
前五条比第一部分那些好得多——它们至少是显式配置,可以被 review、被监控。但它们仍然是分散的、彼此独立的决策,而且没有一处表达了它们真正想说的东西。
比如"重灌每天跑一次",它真正想说的是"这张表的数据可以容忍一天的收敛延迟";“这张表建在 Fluss”,真正想说的是"这张表要秒级新鲜"。我们一直在用实现手段表达意图,而不是直接声明意图。
这就是最后那个缺口:需求变化时——算法侧提速了、某个大屏要求更快、某条链路成本太高要降级——你要做的还是老三样:改代码、改配置、重新对账。第一部分我们说"时效性本该是一个参数",第二部分换掉了地基,但那个参数依然不存在。
二、Materialized Table:把刷新方式从表定义里分离出来
Flink 的 Materialized Table 把这件事翻转过来:你只声明结果应该是什么、以及要多新,剩下的交给框架。
一张表的定义包含两部分——一段描述业务语义的查询,和一个 FRESHNESS 声明。可以先这样锚定:它像数据库里的物化视图,但多了一个"要多新"的维度。物化视图只回答"这张表是什么",Materialized Table 还回答"它需要多新",而后者决定了它怎么被算出来。
框架根据 freshness 决定执行形态:要求秒级或分钟级,跑成一个持续的流式作业做增量刷新;要求小时或天级,退化成周期性的批作业做全量刷新。查询侧看到的始终是同一张表,不感知底下用的是哪一种。
对照第二部分就能看出差别:那里"这张表放 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 秒"意味着什么,而不是等作业跑起来才发现集群不够用。
六、三个阶段的一条线
三个阶段其实在回答同一个问题的三个层次:同一份业务语义,如何在不同时效要求下保持一致。
现状的答案是"各写一遍,人来对齐";第二阶段是"写一遍,存储来对齐";第三阶段是"写一遍,时效是它的一个参数"。
值得注意的是每一步消失的东西:第一步消失的是三套代码和多写入者竞争,第二步消失的是重灌调度、tag 时机、批次号单调这些编排。架构演进的度量单位不是加了什么组件,而是少了多少需要人记住的东西。
结语
这篇文章真正想说的其实不是某个组件好用,而是一个判断标准:
一个架构的复杂度,不在于它有几个组件,而在于有多少约束只存在于人的记忆里。
第一部分那些重放起点、保留期对齐、对账阈值,第二部分剩下的重灌节奏和 tag 时机——每一条被系统接管,架构就真正简单了一点。反过来,加再多组件,只要约束还在文档和某个人的脑子里,复杂度就没有下降。