跳到主要内容
Apache-2.0Rust · Tokio · ArrowCNCF Landscape

一套永不丢失记录的流处理引擎

高性能 Rust 流处理引擎。用声明式 YAML 编写流水线,以 SQL 处理数据,依托预写日志(WAL)持久化、检查点状态,以及内置的控制平面完成机群运维。

./arkflow --config config.yaml —— 单一二进制,无需集群。

config.yaml

为生产级数据流水线而生

从笔记本上的单个二进制,到经过审计的机群发布。

高性能

Rust + Tokio,列式 Apache Arrow 数据模型。算子融合为链——流水线内部没有跨通道跳转。

默认持久化

每条消息在处理前先 fsync 到预写日志。游标只按最高连续确认序列号推进。

可选精确一次

Kafka 事务输出端到端消除重复窗口。检查点状态将作业恢复到一致切面。

用 SQL 处理

DataFusion 驱动的 SQL、窗口函数、Python UDF 与 VRL 变换——另有 Protobuf、Debezium CDC 与 Schema Registry 编解码器。

机群控制平面

Hub/Agent 架构,期望状态语义、调和机制、可审计的发布流程,以及带作业 DAG 编辑器的 Web 控制台。

同一内核,任意拓扑

流与分布式作业编译为同一 JobSpec。你只需声明 输入 → 缓冲 → 处理器 → 输出,其余交给内核。

一份配置,贯穿全链路

接入
KafkaHTTPMQTTNATSPulsar
处理
SQLPythonVRL窗口
投递
SQLKafkaRedisInfluxDBS3 WAL
✓ 处理前先落 WAL 持久化✓ 默认至少一次,可选精确一次✓ 检查点、崩溃、恢复

深入引擎内部

流与分布式作业编译为同一 JobSpec,运行在同一个执行内核上。批次保持持久,SQL 承担重活,输出保持有序——同时 Hub/Agent 控制平面守护整个机群。

ArkFlow 引擎架构数据源接入 ArkFlow 引擎——经输入、WAL、缓冲、处理器、输出,以 Arrow record batch 流转——由 Hub/Agent 控制平面统一调度,并写入下游。控制平面 · ARKFLOW-SERVERWeb 控制台HubAgent期望状态 · 调和数据来源输出目标ARKFLOW 引擎流与作业编译为同一 JobSpec · 统一执行内核输入源读取WALfsync · 重放缓冲窗口 · 连接处理SQL · UDF · VRL输出有序写出错误输出Arrow MessageBatch有界通道 · 背压检查点状态Kafka消费者组MQTTQoS 订阅HTTPREST 推送NATScore · JetStreamMySQL批量 upsertKafka事务Redis热路径缓存InfluxDB指标
悬停任意节点查看说明 —— 图中的光波就是一个流经内核的批次。
流入记录 128,940已处理 128,940丢弃 0

加入 ArkFlow

ArkFlow 是基于 Apache-2.0 的开源项目,已收录进 CNCF Landscape。欢迎贡献代码、提交 Issue 与分享想法。