arkflow
Rust 原生高性能流处理引擎,融合 Arrow 列式存储 + DataFusion SQL 引擎,
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
Rust 原生高性能流处理引擎,融合 Arrow 列式存储 + DataFusion SQL 引擎,
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
ArkFlow 是一款由 Rust 语言编写的高性能流处理引擎,它将 Apache Arrow 列式存储、DataFusion SQL 引擎与 Tokio 异步运行时三者融合,为数据工程师提供一套可插拔的实时 ETL 管道框架——支持 Kafka、MQTT、NATS、Pulsar、Redis 等 20+ 种数据源,内置 SQL 查询、Python UDF、VRL 重映射等多种处理器,同时无缝集成机器学习模型进行流式推理与异常检测。
图 1:ArkFlow 项目 Logo
在大数据领域,流处理并不是一个新概念。从 Apache Storm 的诞生(2011 年),到 Flink 的崛起(2014 年),再到 Kafka Streams 的轻量化探索,行业里已经有不少成熟的流处理方案。然而,这些系统普遍存在几个共同痛点:
第一,部署复杂度高。 Flink 集群需要依赖 Java 生态和独立的任务管理器,对于只想在单机或小规模场景下处理数据的团队而言,学习成本和运维成本都偏高。
第二,AI 集成能力弱。 传统的流处理框架设计初衷是数据搬运和清洗,并没有把机器学习模型的实时推理当作核心场景。当工程师需要在一个数据流中对图像、文本或传感器数据做 AI 推断时,往往需要自己维护一个模型服务,再通过 HTTP 或 Kafka 与流处理框架对接,架构变得臃肿且延迟上升。
第三,资源占用大。 Java 系的流处理框架(JVM)本身就需要占用可观的内存,即使在 CPU 资源充足的环境下,也可能因为 GC 停顿导致尾延迟(tail latency)波动。
ArkFlow 的作者 chenquan 显然注意到了这些问题。项目采用 Rust 语言编写,利用 Rust 零成本抽象和内存安全的特性,在保证极致性能的同时避免了 JVM 的资源开销。更重要的是,它将 Apache Arrow 的列式数据模型(RecordBatch)作为内部统一数据结构,使得 DataFusion SQL 引擎可以直接对流数据进行查询,无需额外的数据转换。
项目已入选 CNCF Cloud Native 云原生技术全景图,在 Streaming & Messaging 分类中占据一席之地,与 Kafka、Flink 等成熟项目同台展示。
ArkFlow 的代码组织为典型的 Cargo Workspace 三 crate 结构:
arkflow-core — 核心抽象层,定义所有组件的 trait 接口(Input、Output、Processor、Buffer、Codec),提供 Engine(引擎主控)和 Stream(单条数据流)两个核心对象。Stream 是最小调度单元,一条 Stream 包含 Input(输入)、Pipeline(处理器链)、Output(输出)三个部分。
arkflow-plugin — 插件实现层,所有具体的数据源/输出/处理器均在此实现。通过 lazy_static + RwLock<HashMap> 的注册机制实现动态插件加载,用户只需要在 YAML 配置中指定类型名称,系统即可在运行时实例化对应插件。
arkflow — 命令行入口 crate,负责解析配置、初始化插件注册表、启动 Engine。
ArkFlow 内部使用 Apache Arrow 的 RecordBatch 作为统一数据容器。RecordBatch 是列式存储格式,同一列的数据在内存中是连续的 CPU 缓存友好布局,非常适合向量化执行。DataFusion SQL 引擎直接以 RecordBatch 为输入,无需序列化/反序列化开销。
每条消息还附带了标准化的元数据字段(以 __meta_ 前缀标识):
| 字段 | 说明 |
|---|---|
__meta_source | 数据源名称 |
__meta_partition | Kafka 分区编号 |
__meta_offset | 消息偏移量 |
__meta_key | 消息 key |
__meta_timestamp | 消息时间戳 |
__meta_ingest_time | 摄入时间 |
__meta_ext | 扩展键值对 |
这些元数据可以在 SQL 查询中直接引用,实现丰富的数据关联。
一条 Stream 的并发执行模型如下:
thread_num 参数控制并发度背压机制使得 ArkFlow 在面对突发流量时能够自动降速,而不是崩溃或丢失数据。
ArkFlow 的依赖图谱清晰反映了它的设计目标:
| 组件层 | 核心技术选型 | 作用 |
|---|---|---|
| 异步运行时 | Tokio | 多线程异步任务调度 |
| 数据格式 | Apache Arrow + DataFusion | 列式存储 + SQL 查询引擎 |
| 消息队列 | rdkafka、async-nats、pulsar、rumqttc | 多协议支持 |
| Python 集成 | PyO3 | Python UDF 处理器 |
| 日志/追踪 | Tracing + Tracing-subscriber | 结构化日志和链路追踪 |
| HTTP 服务 | Axum | 健康检查接口 |
| 缓存/存储 | Redis、RDKafka | 缓冲和持久化 |
| 对象存储 | object_store crate | S3/GCS/Azure/HDFS 统一访问 |
| 流处理语言 | VRL(Vector Remap Language) | 数据重映射和富化 |
特别值得注意的是 DataFusion 的引入——它是 Apache Spark 生态之外最成熟的 Rust 原生 SQL 引擎,支持完整的 SQL 语义(聚合、窗口函数、JOIN),这意味着用户可以用熟悉的 SQL 语法直接操作实时数据流,而无需学习专用 DSL。
ArkFlow 支持 13 种输入类型,覆盖了主流的物联网、消息队列、数据库和文件场景:
6 种内置处理器满足大多数 ETL 场景:
__meta_ 元数据5 种缓冲策略支持复杂事件处理:
虽然没有独立的 ML processor,但通过 Python UDF 处理器,ArkFlow 可以加载任意 Python ML 库(PyTorch、TensorFlow、ONNX Runtime)进行流式推理。用户只需在 YAML 中指定 Python 脚本路径,系统在运行时加载并对每批 RecordBatch 调用推理函数,结果写回 Arrow 列继续下游处理。
从源码构建需要 Rust 1.88+ 环境:
git clone https://github.com/arkflow-rs/arkflow.git
cd arkflow
cargo build --release
./target/release/arkflow --config config.yaml
Docker 方式构建(生产推荐):
cd arkflow
docker build -f docker/Dockerfile -t arkflow .
docker run -v $(pwd)/config.yaml:/app/config.yaml arkflow
一个从 Kafka 消费数据、做 SQL 过滤、输出到 PostgreSQL 的完整配置:
logging:
level: info
streams:
- input:
type: kafka
brokers: ["localhost:9092"]
topics: ["sensor-data"]
consumer_group: arkflow-consumer
pipeline:
thread_num: 4
processors:
- type: json_to_arrow
- type: sql
query: "SELECT * FROM flow WHERE value > 100"
output:
type: sql
url: postgresql://user:pass@localhost:5432/db
table: filtered_data
整个系统没有 Web UI,所有配置通过 YAML 文件驱动。官方提供了 20+ 份示例配置文件(examples/ 目录),覆盖每一种组件组合,是最实用的参考材料。
作为一个相对年轻的项目(v0.5.0),ArkFlow 的生态远不如 Flink 或 Kafka Streams 丰富。没有内置的 Exactly-Once 语义保证(虽然背压机制提供了 At-Least-Once),状态后端目前仅支持内存和临时存储,大规模有状态流处理能力有待验证。
由于完全依赖 YAML 配置运行,当数据处理出错时,排查链路需要阅读日志或启用 debug 级别的 tracing 输出,缺乏可视化的数据预览和断点调试能力。这对习惯 Flink Web UI 的用户来说是一个明显的落差。
虽然项目提供了健康检查端点(/health、/readiness、/liveness),但没有提供官方的 Helm Chart 或 Kustomize 模板,云原生生产部署需要用户自行适配。
官网文档(arkflow-rs.com)提供了中英文双语版本,但部分高级功能(如 Python UDF、VRL 表达式参考)的文档仍以英文为主,中文开发者需要参考 VRL 官方文档补充学习。
ArkFlow 的出现代表了 Rust 生态在数据基础设施领域的一次重要推进。
Rust 语言本身在数据系统中的优势已经在 ClickHouse(存储引擎)、Apache Arrow(核心数据结构)、DataFusion(SQL 引擎)等项目中得到验证。ArkFlow 将这些 Rust 原生组件串联起来,形成了一条从数据采集到处理的完整链路,且全程无需 JVM 依赖——这对于边缘计算、物联网网关等资源受限场景非常有吸引力。
当前项目已有 1283 个 GitHub Stars,20 个贡献者,Star 趋势稳定增长。作为 CNCF Landscape 的收录项目,它正在逐步建立起行业认可度。后续如果能够补充状态管理、Kubernetes Operator 支持和 Web UI 监控界面,将有望在中小规模的实时数据场景中占据一席之地。

图 2:ArkFlow 微信交流群(欢迎加入讨论)
| 属性 | 值 |
|---|---|
| 编程语言 | Rust |
| 最低 Rust 版本 | 1.88 |
| 许可证 | Apache 2.0 |
| 官方文档 | arkflow-rs.com |
| 数据模型 | Apache Arrow RecordBatch |
| SQL 引擎 | Apache DataFusion |
| 异步运行时 | Tokio |
| CI/CD | GitHub Actions(Rust CI badge) |
| 入库 CNCF | 是(Cloud Native Landscape) |
| Docker 支持 | 多阶段 Dockerfile(生产级) |
| Web UI | 无(YAML 驱动) |