Streamling 的三个架构决策

14 阅读 2473 字 · 约 9 分钟

Streamling 有三个架构决策让我停下来想了很久。这些选择影响的是系统行为本身。

我做过 Rust + Actor 模型的流处理引擎,对这些问题有自己的偏好。所以每个决策我都会从"如果是我会怎么选"的角度来聊。

Checkpoint 跟着数据走

Streamling 的 checkpoint marker 不走独立通道,而是搭载在 RecordBatch 的 schema metadata 上,跟着数据一起流过整个拓扑。

每条数据批次自带一个 schema,schema 里有一段 metadata 记录"这是第几个 checkpoint epoch"。数据流过算子的时候,checkpoint 信息跟着数据一起到。到 sink 的时候,sink 确认这个 epoch 的数据已经写入,发一个 ack 回协调器。

如果是我,我不会这么做。

它依赖一个隐式契约:所有算子都必须保留 schema metadata。DataFusion 原生算子不保证这一点。DataFusion 是为批查询设计的,一个 FilterExec 执行完可能重建 schema,metadata 就丢了。每升一个 DataFusion 版本,都可能引入不保留 metadata 的新算子。这是一个脆弱的依赖。

我会用独立的消息通道传 checkpoint marker,跟数据流解耦。数据走数据的通道,checkpoint marker 走 checkpoint 的通道。清晰,但代价是多一个通道要管理。可能出现 marker 到了但数据还没到的对齐问题。两边的 marker 和数据要匹配,不能错位。

他们为什么选了这个方式?有一个好处:marker 跟数据严格同步,不存在对齐问题。数据到哪,marker 就到哪。不会出现 marker 领先数据或落后数据的情况。对于 exactly-once 语义来说,这个保证很重要。

紧凑,但脆。

用优化器规则注入流处理语义

上面那个问题有个续集。DataFusion 原生算子会丢掉 schema metadata,那怎么办?

Streamling 注册了自定义的物理优化器规则。在执行计划阶段,把 DataFusion 原生的 FilterExec 替换成自己写的 StreamingFilterExec,ProjectionExec 替换成 StreamingProjectionExec。替换后的算子会保留 checkpoint marker metadata。

DataFusion 是为批查询设计的,没有流处理的概念。Streamling 没有 fork DataFusion,而是用框架提供的扩展点注入流处理语义。

这个做法我认同。

fork 一个框架,你拿走了源码控制权,也承担了所有后续维护成本。上游修了一个 bug,你得手动 cherry-pick。上游加了一个新功能,你得评估要不要合进来。版本越拖越远,最后变成一个没人维护的分叉。

跟上游版本走,你被动适配。上游改了接口,你就得改。上游加了一个不保留 metadata 的算子,你的 checkpoint 就断了。

用优化器规则注入行为,是第三条路。你用框架的扩展点,在框架之上加一层自己的逻辑。框架升级的时候,只要扩展点接口没变,你的逻辑就不受影响。这是扩展点设计的正确用法。

不过我会加一个改进。把改写规则做成可配置的。用户在升级 DataFusion 的时候,可以临时关掉改写规则,先验证新版本兼容性,再打开。相当于给升级过程加了一个安全网。目前看代码里没有这个开关。小项目无所谓,但如果未来用户多了,DataFusion 大版本升级会是一个风险点。

插件拿到宿主的 Tokio Runtime

Streamling 的插件通过 FFI 加载。加载之后,插件拿到的是宿主的 Tokio runtime handle。插件可以 spawn 异步任务、sleep、block,不需要自己起 runtime。

FFI 场景下异步 runtime 的所有权是个经典痛点。插件是一个动态链接库,它跟宿主跑在同一个进程里。如果插件自己起一个 Tokio runtime,就有两个 runtime 在同一个进程里抢资源。tokio 的 runtime 不是设计来共存的,两个 runtime 互相 spawn 任务会出问题。

如果插件不起 runtime,它就没法做异步 IO。一个需要调 HTTP API 的 transform 插件,没有 async 就只能阻塞,性能直接垮掉。

Streamling 的解法:宿主把自己的 runtime handle 传给插件。插件用宿主的 runtime 来 spawn 任务。runtime 只有一个,不存在冲突。插件想做异步 IO,随时可以。

简洁。但我会在安全边界上加一层。

插件拿到的是宿主的 runtime handle,可以 spawn 任意任务。没有限制。如果插件 panic 了,panic 会传播到宿主。如果插件 spawn 了一个无限循环的任务,宿主的线程池里就多了一个永远跑不完的任务。在 Goldsky 的场景下插件是半可信的:用户自建 pipeline,插件是用户自己写的,不是恶意攻击者。这个取舍目前能接受。

但如果未来开放给不信任的第三方插件,这就是个安全隐患。我会给插件一个受限的 runtime wrapper。限制并发任务数,加 panic 隔离。插件 spawn 的任务如果 panic,只影响插件自己,不波及宿主。相当于在 runtime 层面加一个沙箱。

FFI 的安全边界是一个真实的问题。Rust 的 unsafe 块可以隔离内存安全,但 panic 传播和资源占用不在 unsafe 的管辖范围内。需要工程手段来解决。

三个决策的共同线索

这三个决策有一个共同特征:都是在"简洁"和"稳健"之间选了简洁。

Checkpoint 搭载在数据上,省了一个通道,但依赖隐式契约。优化器规则注入,省了 fork 的成本,但依赖扩展点稳定。插件共享宿主 runtime,省了一个 runtime 的复杂度,但放弃了隔离边界。

每个决策都解决了真实的问题,在当前场景下都 work。但简洁的代价是边界条件更脆。随着项目规模和用户群体增长,这些边界条件会被压力测试。

这些是技术债。每个项目都有。关键是团队知不知道这些债在哪里,什么时候该还。从代码注释来看,他们知道。