流式数据库深度实战:从增量视图维护、Differential Dataflow 的 Arrangement 索引到流式 SQL 一致性模型的工程全解

在过去五年里,"实时"从一个差异化卖点变成了基础设施的及格线。风控需要在 200 毫秒内判定一笔交易,推荐系统需要在用户完成下一次点击前刷新特征,可观测性平台需要在指标越界的瞬间给出归因。业界最初的解法是"流计算引擎 + 外部存储":用 Flink 或 Spark Streaming 读 Kafka,算完写回 Redis、ClickHouse 或 Hudi,再让上层服务去查。

这套架构能用,但它把复杂度转嫁给了业务方:状态语义要自己对齐、维表关联要自己做缓存与刷新、结果表的更新时序要自己保证、口径一旦变化就要重跑全量回溯。流式数据库(Streaming Database)试图从根上换一种做法——让用户只写 SQL,让系统自己维护物化视图的增量一致性。RisingWave、Materialize、DeltaStream 是这一路线的代表,它们背后的共同理论支柱是增量视图维护(Incremental View Maintenance, IVM)与 Differential Dataflow。

这篇文章聚焦三件事:增量视图维护到底在维护什么;Differential Dataflow 为什么能用一套"差分集合"抽象同时表达迭代、join 与窗口;以及在真实工程里,流式 SQL 的一致性边界究竟在哪里。


一、从批处理到 IVM:把"重算"变成"求导"

先看一个最朴素的聚合视图:

CREATE MATERIALIZED VIEW order_stats AS
SELECT seller_id,
       count(*)             AS order_cnt,
       sum(amount)          AS gmv,
       avg(amount)          AS avg_amount
FROM orders
GROUP BY seller_id;

批处理的语义是:每次数据变化,重新扫描全表重算。流式数据库的语义是:给定输入集合的差分 ΔR,计算输出集合的差分 ΔV,使得 V ⊕ ΔV = Q(R ⊕ ΔR)。这本质上是把查询 Q 看成函数,对输入求导。

对于纯累加的算子(count、sum、min/max on append-only),这很容易:

  • count(*) 的增量就是"新增行数"
  • sum(x) 的增量就是"新增行的 x 之和"

但一旦出现需要撤回(retraction)的算子,事情就复杂了:

-- 分组取 Top-3 卖家,按 GMV 排序
CREATE MATERIALIZED VIEW top_sellers AS
SELECT seller_id, gmv,
       rank() OVER (ORDER BY gmv DESC) AS rk
FROM order_stats
WHERE rk <= 3;

一条新订单进入后,order_stats 中某个 seller 的 gmv 上升,可能挤掉原本第 3 名的卖家。系统必须主动撤回旧的第 3 名行(发出一个 -1 的差分),再发出新的第 3 名行(+1)。如果撤回链路在某个算子里断了,视图就会静默地产生错误结果——这是流式数据库最危险的一类 bug,因为它不报错,只是"看起来一直很新但数字不对"。

这也是为什么流式数据库普遍采用差分数据流(Differential Dataflow)而非简单的"消息流"作为内部表示。


二、Differential Dataflow:用 (data, time, diff) 三元组统一一切

Differential Dataflow(源自 Frank McSherry 的 Naiad 工作)把集合的演化表示为多重集(multiset)上的差分集合:

集合 R 在时间 t 的内容 = Σ over all (data, time, diff) where time <= t  of  diff * data

每条记录携带三个字段:

  • data:记录本身
  • time:逻辑时间戳(在 Naiad/Timely Dataflow 中是一个偏序的 Antichain 元素,而非标量)
  • diff:多重性变化量(通常是 +1 / -1,聚合场景可以是任意整数)

用 Rust 写一个最小的差分 join,语义非常清晰:

use differential_dataflow::input::Input;
use timely::dataflow::operators::Inspect;

fn main() {
    timely::execute_from_args(std::env::args(), |worker| {
        let mut orders   = worker.dataflow::<usize, _, _>(|scope| { /* ... */ });
        // 简化示意:两个输入集合做等值 join,输出差分
        worker.dataflow::<usize, _, _>(|scope| {
            let (o_handle, orders)   = scope.new_collection();
            let (u_handle, users)    = scope.new_collection();

            // key -> (user_id, amount)  与  key -> (user_name)
            let orders_by_user = orders.map(|(uid, amount, item)| (uid, (amount, item)));
            let users_by_key   = users.map(|(uid, name)| (uid, name));

            // differential join:两侧任一边的差分都会驱动另一侧的增量查找
            users_by_key
                .join(&orders_by_user)
                .map(|(_uid, (name, (amount, item)))| (name, item, amount))
                .inspect(|(x, t, d)| println!("{:?} @ {:?} diff={}", x, t, d));

            // 插入一条用户 + 两条订单,随后撤回一条订单
            u_handle.insert(1, ("alice".to_string(),));
            o_handle.insert((1, 100, "keyboard".to_string()));
            o_handle.insert((1, 30,  "mouse".to_string()));
            o_handle.advance_to(1);
            o_handle.remove((1, 30, "mouse".to_string())); // 发出 diff = -1
            o_handle.advance_to(2);
        });
    }).unwrap();
}

关键在于 join 的实现:它不是"流式窗口 join",而是在索引(Arrangement)上做增量查找。当左集合收到 (k, v, t, +1),系统去右集合的 arrangement 里查 k,对命中的每个 w 输出 (k, (v,w), t, +1)。撤回同理。因此 join 的复杂度是 O(|Δ| · log N),而不是 O(N)。

Arrangement:把集合变成可增量查询的索引

Arrangement 是差分集合的物化索引版本——它把 (data, time, diff) 按 key 组织成有序结构(通常是类 LSM 的 trace 层,磁盘上是不可变批次 + 内存层做合并)。有了它:

  • join 变成索引查找
  • group / reduce 变成按 key 的局部聚合
  • count/sum 的撤回变成"在 key 的累积值上加 diff"

Materialize 的 compute 层正是如此:每个 dataflow 算子维护自己的 arrangement,通过 spine / trace 结构支持按时间做历史查询。RisingWave 采取了稍微不同的路线——它把状态存成 LSM-Tree 里的 KV(与流计算引擎更接近),通过 StateStore 抽象支撑 Hummock(对象存储上的 LSM)做状态持久化。两者差异可以概括为:

维度MaterializeRisingWave
计算内核Differential Dataflow(Timely)自研 Streaming Engine(Rust)
状态表示in-memory arrangement trace + persistLSM StateStore(Hummock,S3 原生)
强项任意嵌套迭代、递归查询、强一致读大状态、云原生弹性、存算分离成本
时间模型偏序 antichain 前沿epoch/barrier 单调递增水位

三、流式 SQL 的一致性边界:这是工程上最容易踩的坑

"只写 SQL"听起来很美,但流式 SQL 的正确性边界比批处理严格得多。几个必须讲清楚的点:

3.1 撤回(retraction)必须端到端保序

考虑一个链式视图:源表 → 聚合 → Top-N → 下游 sink。Top-N 产生的 -1 与 +1 如果乱序到达 sink,下游会短暂地看到"重复的第 3 名"。流式数据库的做法是给每条变更打上 epoch(或 frontier),并保证同一 key 的变更严格按 epoch 递增。RisingWave 的 barrier 会在整个 DAG 上做对齐,Materialize 则用 frontier 的偏序推进来表达"某个时间点之前的数据已经完整"。

消费端必须正确处理这两种语义。以 RisingWave 的 SUBSCRIBE 为例:

-- 订阅物化视图的变更流,而不是查询快照
SUBSCRIBE top_sellers WITH (snapshot = 'process');

返回的每一行带 op(Insert/Delete/UpdateInsert/UpdateDelete)与 rw_timestamp。消费端代码示例:

import psycopg2

conn = psycopg2.connect("host=risingwave dbname=dev user=root")
cur  = conn.cursor()
cur.execute("SUBSCRIBE top_sellers WITH (snapshot = 'process');")

seen_epoch = None
for (op, epoch, seller_id, gmv, rk) in cur:
    if epoch != seen_epoch:
        # 一个 epoch 的变更原子提交,这里是做"批次级"对外可见的边界
        flush_downstream()
        seen_epoch = epoch
    if op in ("Insert", "UpdateInsert"):
        upsert(seller_id, gmv, rk)
    elif op in ("Delete", "UpdateDelete"):
        # 必须处理撤回,否则 Top-N 结果会无限膨胀
        delete(seller_id)

工程观点:绝大多数"流式数据库不准"的抱怨,最后都定位到消费端忽略了 Delete/UpdateDelete。撤回不是可选项,它是一致性模型的一部分。

3.2 维表关联(Temporal Join)的时态语义

流式事实表关联缓慢变化维表时,正确的问题是"在事实发生的那个时刻,维度是什么",而不是"维度现在是什么":

-- RisingWave:按事件时间做时态关联
CREATE MATERIALIZED VIEW enriched AS
SELECT o.order_id,
       o.amount,
       u.user_tier,          -- 下单那一刻的会员等级
       o.event_time
FROM orders o
LEFT JOIN users FOR SYSTEM_TIME AS OF o.event_time AS u
  ON o.user_id = u.user_id;

这要求系统保存维度的历史版本(本质上又是一张 arrangement,key = 主键,time = 版本时间)。代价是状态膨胀;收益是回溯计算可复现。若不写 FOR SYSTEM_TIME,你得到的是"当前维度"关联,结果随维表更新而历史漂移——这在财务对账场景是致命的。

3.3 窗口与水位线:迟到数据的三重处理

CREATE MATERIALIZED VIEW hourly_gmv AS
SELECT window_start, seller_id, sum(amount) AS gmv
FROM TUMBLE(orders, event_time, INTERVAL '1' HOUR)
GROUP BY window_start, seller_id;

必须显式决定迟到数据的策略,否则水位线推进后到达的数据会被静默丢弃:

-- 允许 5 分钟迟到,期间窗口结果会被反复撤回与重发
SELECT window_start, seller_id, sum(amount)
FROM TUMBLE(orders, event_time, INTERVAL '1' HOUR, INTERVAL '5' MINUTES)
GROUP BY window_start, seller_id;

经验值:迟到容忍期越长,状态保有量越大,撤回放大越严重。一个允许 24 小时迟到的 tumble 窗口,其状态量约等于 24 小时的输入量——这在对象存储上很便宜,但在内存 arrangement 上很贵。选型时要明确问一句:我的撤回窗口有多大?


四、什么时候不该用流式数据库

我见过不少团队把流式数据库当成"更快的数仓",然后在三周后痛苦地回退。三条硬边界:

  1. 需要复杂多路回溯重算的指标。全量回溯在 IVM 下等价于"从零重放",成本可能高于一次批处理。若你的核心诉求是"每天重算过去两年",用批。
  2. 结果需要强一致多表事务。物化视图之间的一致性是"同一 epoch 内一致",跨 sink(比如同时写 Postgres 和 Kafka)没有分布式事务保证。需要跨系统原子性就得上两阶段提交或事务消息,别指望数据库替你解决。
  3. 状态规模远超热数据且撤回窗口极长。此时内存 arrangement 的持有成本会失控;应该退回到"Flink/Spark + 外部 KV"的架构,或者像 RisingWave 那样把状态下沉到对象存储并接受更高的读延迟。

反过来,以下场景是流式数据库的主场:

  • 指标口径稳定、但数据持续到达(实时监控、风控特征、计费)
  • 需要"查到即最新"的交互式延迟(亚秒级),而不是分钟级微批
  • 团队有 SQL 能力但没有 JVM 流计算工程能力

五、一条可执行的落地路径

  1. 先做烟囱式试点:选一个现有批处理链路(如小时级 GMV 聚合),用同一个 SQL 建物化视图,双跑两周对账。对账要在 epoch 边界上做,而不是随机抽样。
  2. 把撤回当成一等公民:sink 端统一实现 upsert + delete 语义,写自动化测试注入撤回事件。
  3. 给每个视图标注数据契约:输入源的乱序程度、迟到容忍、撤回窗口、状态大小上限。这些参数比 SQL 本身更容易引发生产事故。
  4. 监控 frontier / epoch 推进:前沿停滞通常意味着上游某个 partition 卡住或某条 barrier 无法对齐,是流式系统最典型的"静默故障"。
  5. 成本模型提前算:状态字节数 × 副本 × 保留时长,加上撤回放大系数(实测通常在 1.5–3 倍),再决定是否上云托管。

六、小结

流式数据库不是一个"更快的 Kafka 消费者",它把数据库里最成熟的那套东西——物化视图、查询优化、事务性状态——搬到了持续到达的数据上。理解它的关键不在于记住多少 SQL 语法,而在于建立三个心智模型:

  • 差分心智模型:一切更新都是 (data, time, diff) 三元组,撤回与插入同等重要。
  • 索引心智模型:join 与聚合之所以能增量,是因为 arrangement 把集合变成了可增量查询的索引。
  • 时间心智模型:epoch / frontier 定义了"什么是完整的",所有可见性语义都建立在它之上。

把这三件事想清楚,你会发现流式 SQL 并不是黑魔法:它只是把"重算"换成了"求导",而求导的代价,是必须诚实地面对撤回、乱序与状态。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部