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 モデルのロードやトランザクションシンク以外の正確に一度の保証については詳細が説明されていません。

最適なユーザー クラウドネイティブで高スループットのストリーム処理エンジンが必要で、ストリーミングデータ上で機械学習モデルを実行し、リアルタイム推論、異常検出、異種データソース間での複雜イベント処理を行いたい開発者およびオペレーター。

関連

  • プロジェクト
  • プロジェクト
  • プロジェクト
  • プロジェクト
  • プロジェクト