arkflow-rs/arkflow
High performance Rust stream processing engine seamlessly integrates AI capabilities, providing powerful real-time data processing and intelligent analysis.
ArkFlow
摘要 高性能的 Rust 流处理引擎,集成 AI 能力用于实时数据处理、推理、异常检测和复杂事件处理。
目的 提供一个持久且可扩展的平台,从多种来源摄取数据,使用 SQL、Python UDF、JSON/Protobuf 转换等进行处理,并将结果写入各种接收端,同时启用机器学习模型的加载与执行,以实现智能流分析。
工作原理
- 基于 Rust 和 Tokio 异步运行时,实现低延迟、高吞吐量处理。
- 用户在 YAML 配置文件中定义流:一个 输入(Kafka、MQTT、HTTP、文件、生成、SQL 等)、一个带有处理器的 管道(json_to_arrow、SQL、Python UDF、VRL、批处理、Protobuf)、一个 输出、可选的 错误输出、以及可选的 缓冲区(内存、会话/滑动/翻滚窗口)。
- 事件以异步方式通过管道流动;通过每个流的预写日志(WAL)实现持久性,提供至少一次交付,且在事务型接收端可获得恰好一次保证。
- 可选的 Hub 和 Web 控制台形成控制平面,用于观察、配置和操作多个 ArkFlow 节点作为舰队。
主要特性
- 高性能:Rust + Tokio。
- 持久交付:基于 WAL 的至少一次交付,事务型接收端可选恰好一次。
- 多数据源:Kafka、MQTT、HTTP、文件(CSV/JSON/Parquet/Avro/Arrow)、SQL、NATS、Pulsar、Redis、WebSocket、Modbus、内存、生成,以及可组合的多输入。
- 强大处理:SQL 查询、Python UDF、JSON/Protobuf 转换、VRL、批处理。
- 流编解码:JSON、Protobuf、Debezium CDC 信封、Confluent Schema Registry 线路格式。
- 控制平面:Hub + Web 控制台,用于舰队管理。
- 可扩展:模块化设计,可添加新的输入、缓冲区、输出、处理器组件。
- 缓冲区:内存、会话窗口、滑动窗口、翻滚窗口,用于背压和窗口化。
局限性
- README 未具体说明任何限制;关于 AI/ML 模型加载或事务型接收端以外的恰好一次保证的细节未详细阐述。
最适合对象 需要云原生、高吞吐量流处理引擎的开发者和运维人员,该引擎能够在流数据上运行机器学习模型,以实现实时推理、异常检测和跨异构数据源的复杂事件处理。
相关
- 项目
- 项目
- 项目
- 项目
- 项目