Streamling:不做分布式的流处理引擎
我之前做过一个基于 Rust + Actor 模型的生产级流数据处理引擎。消息传递,有状态计算,那套工具箱我很熟。
看到 Streamling 的时候,我以为又是同类项目。Rust 写的,流处理,能有什么不同?
但看了项目介绍,我发现这个开源项目已经在银行、对冲基金、预测市场的生产环境跑着。读完核心源码后,我改了看法。它做了一个跟主流流处理引擎不同的架构选择。
Goldsky 为什么造这个项目
Goldsky 做区块链数据索引。他们的产品 Turbo 让用户自建数据管道:从链上读取原始事件,解码变换,写入查询存储。客户是银行、对冲基金、预测市场,对数据准确性和实时性要求很高。
为什么是上千个用户而不是几个大企业?区块链数据是公开的,但原始的事件日志和状态变更需要解码成结构化数据才能查询。每个 DApp、DeFi 协议、分析团队需要的数据切片不同。一个协议关注特定合约的事件,一个分析团队要特定 token 的转账记录。需求碎片化,但加起来量很大。
传统金融数据是私有的、集中的,几个大机构各自建数据团队。区块链数据是公开的、去中心化的,参与者大量且分散。Turbo 是自助式 SaaS,用户在 web 界面配置管道,不需要运维团队。门槛低意味着中小型用户也能进来。银行、对冲基金、预测市场是其中对数据质量要求最高的客户,但整个生态的参与者远不止这几类。上千条 pipeline 对应的是整个区块链生态的索引需求。
问题在于规模。不是一个团队跑一条大管道,而是上千个用户各自跑自己的小管道。每条管道的 source、transform、sink 都不同。用 Flink 跑这个场景,每条管道要起一个 Job,JVM 基础开销几百 MB 起步,上千条管道就是几百 GB 内存。这还没算数据。运维上,每条管道要部署、监控、独立升级,Flink 的集群管理在这种"大量小任务"模式下反而成了负担。
Goldsky 需要一个引擎:每条管道内存开销小到可以忽略,启动快到秒级,用户能自定义算子且不需要重编译引擎,声明式配置让非专业用户也能上手。市面上没有现成方案满足这些条件。所以他们自己写了 Streamling。
Streamling 是什么
Streamling 是一个用 Rust 写的列式流处理运行时,处理实时和历史数据。你用 YAML 声明数据管道:从哪里读(Kafka、文件)、怎么变换(SQL、WASM、HTTP 调用)、写到哪里(Postgres、ClickHouse、Kafka、webhook)。单二进制启动,不需要集群。
核心特性:
- 列式数据平面:Arrow RecordBatch 在算子间零拷贝传递,不序列化
- 插件系统:FFI 动态加载,独立 Rust crate,不重编译引擎
- Checkpoint 持久化:at-least-once 语义,状态存 Postgres
- 批流一体:同一套管道逻辑跑实时和批处理
- 内置连接器:Kafka、Postgres、ClickHouse、文件、webhook 开箱即用
适合的场景:实时 ETL(Kafka 数据变换写入分析存储)、CDC 数据同步、事件流索引、批处理任务。不适合需要跨分区 join、窗口聚合、分布式协调 checkpoint 的场景。
它是 Goldsky Turbo 的底层引擎。Goldsky Turbo 在生产环境跑了上千条用户自建的数据管道。
项目地址:github.com/Angryshark128/streamling
架构核心:列式数据平面
大部分流处理引擎的核心抽象是算子。算子之间通过消息传递数据,每个算子维护自己的状态。Actor 模型、Flink 的算子链、Kafka Streams 的拓扑,走的都是这条路。
Actor 模型擅长细粒度的有状态计算。每个 actor 维护自己的状态,消息传递天然支持分布式扩展。但数据在 actor 之间传递需要序列化,列式处理的优势发挥不出来。
Streamling 走了另一条路。它的核心抽象是 Arrow RecordBatch。所有算子之间流动的数据都是 Arrow 的列式内存格式。SQL 变换交给 DataFusion 执行,DataFusion 原生处理 Arrow。Source 解码后产出 Arrow,Sink 消费 Arrow 写入目标存储。整个管道里没有序列化开销。
列式布局带来三个直接好处:CPU cache 友好,内存紧凑,SIMD 指令可以直接作用在连续内存上。
如果是我选,Arrow 数据平面我也会选。但我会在算子级状态上犹豫。Actor 模型里每个 actor 可以维护自己的状态,做窗口聚合、会话计算都很自然。Streamling 的算子近乎无状态,状态只存在于 checkpoint 和 sink 里。这是取舍。代价是窗口聚合这类场景做不了,用户得自己想办法。
README 提到一个基准测试。Streamling 在 Kafka 场景下每 GB 内存的吞吐大约是 Flink 的 40 倍。一个客户写百万行 batch 到 ClickHouse,内存从 31GB 降到 4GB。40 倍不是魔法。列式内存布局加 Rust 无 GC,架构选型决定的天花板就在那里。
不做什么,比做什么更重要
Streamling 明确说不做分布式有状态处理。跨分区 join、窗口聚合、跨节点协调 checkpoint,都不碰。Flink 的复杂度大部分来自分布式协调:JobManager/TaskManager 架构、barrier 对齐、RocksDB state backend、savepoint 管理。砍掉这些,整个系统简单了一个量级。单二进制启动,checkpoint 存 Postgres,不需要 ZooKeeper,不需要 K8s 协调层。
单节点是我最犹豫的地方。没有故障转移,节点挂了管道就断了。但换个角度想,他们的场景是上千条独立的小管道,每条管道数据量不大。与其做一个复杂的分布式系统服务少数大用户,不如做一个简单的单节点系统服务大量小用户。管道级别的高可用可以用进程重启加 checkpoint 恢复来解决,不一定要分布式。这个判断我认为是对的,但前提是单管道吞吐足够高,不需要水平扩展。目前看 Arrow 加 Rust 的性能余量撑得住。
如果你需要窗口聚合和跨分区 join,Flink 仍然是正确选择。Streamling 不试图覆盖所有场景。
发展空间
引擎本身跟区块链没关系。任何"Kafka 读数据、变换、写存储"的场景都是射程范围。CDC 数据同步、实时数仓入仓、日志指标管道、微服务间事件同步,需求跟区块链索引相似:单管道吞吐要求高,但不需要跨分区 join 和窗口聚合。
我最看好的方向是嵌入式流处理引擎。SQLite 没有替代 PostgreSQL,它创造了一个新品类。大量数据管道场景不需要 Flink 集群,但需要一个可靠的流处理引擎嵌入到产品里。SaaS 产品想让用户自建数据管道,IoT 边缘设备需要本地实时处理。Flink 太重,自己写太贵。单二进制、低内存、秒级启动、声明式配置、插件扩展,这套组合目前没有对手。
AI 数据管道也是一个方向。实时特征工程现在大量用 Flink 做,运维成本让很多团队望而却步。Arrow 已经是 ML 生态的通用格式,Streamling 的数据平面原生就是 Arrow,跟 ML 工具链零摩擦对接。
能不能兑现取决于 Goldsky 愿不愿意投入做社区建设,以及有没有第二批用户拿它做跟区块链无关的产品。
为什么选这个项目做深度研究
四个理由。
代码量可控。 12 万行 Rust,10 个 crate。一个人可以读完核心模块。Flink 的代码库是百万行 Java,想深入理解每一层的成本差一个量级。
架构选择非平凡。 Arrow 数据平面、checkpoint 搭载在数据上、物理优化器改写、FFI 插件系统。每一个都值得单独拆解,影响系统行为。
生产验证。 Goldsky Turbo 的客户规模说明这个架构在真实负载下站住了。
可贡献。 82 个 commit,12 个贡献者。项目活跃但规模不大。深入了解源码的人可以找到有价值的贡献点。
接下来
这个系列会逐个拆解 Streamling 的核心模块。每篇深入源码,分析设计决策的 why。计划覆盖:三个架构决策的详细分析、Arrow 数据平面的实际流转、checkpoint 与 exactly-once 的完整机制、FFI 插件系统的设计细节、声明式拓扑引擎的构建过程。
寒蝉 Hancic


