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 與網頁主控台構成控制平面,用於觀察、設定及操作多個 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 + 網頁主控台,用於艦隊管理。
- 可擴充:模組化設計,可新增輸入、緩衝區、輸出、處理器元件。
- 緩衝區:記憶體、會話視窗、滑動視窗、翻轉視窗,用於背壓與視窗化。
限制
- README 未特別說明任何限制;關於 AI/ML 模型載入或交易型接收端以外的恰好一次保證的細節未闡述。
最適合對象 需要雲端原生、高吞吐量串流處理引擎的開發者與操作人員,該引擎能在串流資料上執行機器學習模型,以進行即時推論、異常偵測及跨異質資料來源的複雜事件處理。
相關
- 專案
- 專案
- 專案
- 專案
- 專案