arkflow-rs/arkflow

High performance Rust stream processing engine seamlessly integrates AI capabilities, providing powerful real-time data processing and intelligent analysis.

ArkFlow

요약 실시간 데이터 처리, 추론, 이상 탐지 및 복잡한 이벤트 처리를 위한 AI 기능을 통합한 고성능 Rust 스트림 처리 엔진입니다.

목적 다양한 소스에서 데이터를 수집하고, SQL, Python UDF, JSON/Protobuf 변환 등으로 처리한 결과를 다양한 싱크에 기록하면서, 머신러닝 모델의 로드 및 실행을 가능하게 하여 지능형 스트리밍 분석을 지원하는 내구성 있고 확장 가능한 플랫폼을 제공합니다.

작동 방식

  • Rust와 Tokio 비동기 런타임을 기반으로 구축되어 저지연, 고처리량을 달성합니다.
  • 사용자는 YAML 설정에서 스트림을 정의합니다: 입력(Kafka, MQTT, HTTP, 파일, 생성, SQL 등), 프로세서(json_to_arrow, SQL, Python UDF, VRL, 배치, Protobuf)를 포함한 파이프라인, 출력, 선택적 에러 출력, 선택적 버퍼(메모리, 세션/슬라이딩/텀블링 윈도우).
  • 이벤트는 파이프라인을 통해 비동기식으로 흐르며, 각 스트림별 쓰기-앞 로그(WAL)를 통해 내구성을 확보하고 적어도 한 번 전달을 보장합니다. 트랜잭션 싱크에 대해서는 정확히 한 번 전달도 이용 가능합니다.
  • 선택적 허브와 웹 콘솔이 제어 평면을 형성하여 여러 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 와이어 포맷.
  • 제어 평면: 허브 + 웹 콘솔을 통한 함대 관리.
  • 확장 가능: 모듈식 설계로 새로운 입력, 버퍼, 출력, 프로세서 구성 요소 추가 가능.
  • 버퍼: 메모리, 세션 윈도우, 슬라이딩 윈도우, 텀블링 윈도우(백프레셔 및 윈도우 처리용).

제한 사항

  • README에서는 특정 제한 사항을 명시하지 않으며, AI/ML 모델 로드 또는 트랜잭션 싱크를 넘어선 정확히 한 번 보장에 대한詳細は 설명되지 않았습니다.

최적의 사용자 클라우드 네이티브, 고처리량 스트림 처리 엔진이 필요하며, 스트리밍 데이터에서 머신러닝 모델을 실행하여 실시간 추론, 이상 탐지 및 이질적인 데이터 소스 전반에 걸친 복잡한 이벤트 처리를 수행해야 하는 개발자 및 운영자.

관련

  • 프로젝트
  • 프로젝트
  • 프로젝트
  • 프로젝트
  • 프로젝트