전체 그래프
Papers

Pathway

papersstreamingpathway

상위: Papers

원문 링크

  • doi: (미등록)
  • arxiv: (미등록)

Pathway

1. Architecture

  • Python (API)
    • 선언적 Table API 제공: SQL처럼 join, groupby, window 같은 연산을 직관적으로 정의 가능.
    • 개발자는 “무엇을 할지(what)“만 정의 → 실행 방식(how)은 엔진이 자동 최적화.
    • Python 생태계(NumPy, PyTorch, Transformers 등)와 쉽게 연동 가능 → 커스텀 UDF(User Defined Function) 삽입 가능.
  • Rust (Engine)
    • Event-driven 실행: 이벤트 발생 시 처리(step), 이벤트 없으면 idle(park).
    • Incremental computation: 전체를 재계산하지 않고 변경된 부분만 반영(논문에서 강조).
    • Mini-batch + latency bound: 사용자가 허용 지연 시간(latency bound)을 지정하면 엔진이 자동으로 배치 크기를 조정 → 처리량/지연 간 Pareto 최적 달성.
  • 배치/스트리밍 통합: 동일한 코드/엔진으로 batch와 stream을 모두 실행 가능 (Lambda 아키텍처 불필요).
  • Iterative/graph 연산 지원: 스트리밍 환경에서도 PageRank, Connected Components 같은 반복 연산을 수행 가능 → Spark/Flink 대비 차별점.

2. Pathway Pipeline

2.1 변경 감지 (Detection)

  • refresh_interval 주기로 폴더/DB/API를 확인 → 변경사항 발생 시 이벤트 생성
  • CDC(Change Data Capture) 커넥터도 예: PostgreSQL Debezium, Kafka 토픽).
  • 변경 유형: Insertion / Retraction (취소 후 새로운 버전 삽입).

2.2 데이터 처리 (Processing)

  • Parsing: 다양한 포맷 지원 (PDF, TXT, PPTX, HTML, 이메일 등).
  • Chunking: 토큰 단위 분할, 슬라이딩 윈도우 등 다양한 분할 전략 선택 가능.
  • Embedding:
    • SentenceTransformer 외에도 OpenAI, Cohere, HuggingFace 등 다양한 임베더 연동 가능
    • GPU 가속 지원 → 대규모 실시간 임베딩 처리.

2.3 인덱스 업데이트 (Index Update)

  • In-memory vector store (기본)
  • USearch, Tantivy 기반 하이브리드 인덱스 지원 (semantic + keyword).
  • 증분 업데이트 방식이라 새 데이터 추가/삭제가 바로 반영됨.

2.4 RAG 활용 (Ready for RAG)

  • 쿼리 시 최신 상태 인덱스를 기반으로 관련 chunk 검색.
  • Pathway의 RAG 템플릿 사용 시 → 실시간 문서 갱신 반영 + 프롬프트 생성 자동화.

3. Fault Tolerance

  • State snapshotting: 파이프라인 상태(테이블, 벡터스토어)를 지속적으로 저장.
  • Exactly-once semantics: 동일 이벤트가 중복 처리되지 않도록 보장.
  • 장애 발생 시 snapshot에서 복원 후 이어서 실행.
  • Kafka, DB, Filesystem과 결합 시에도 유실 없는 스트리밍 처리 가능.

4. Extra

  • 튜너블 지연/처리량 트레이드오프: Flink는 설정이 복잡하고 Spark는 microbatch 고정인데, Pathway는 “허용 지연시간”을 하나만 주면 자동 최적화.
  • 멀티테넌시와 확장성: K8s 배포 지원, 워커 단위 확장으로 수평 확장 가능.
  • 실시간 ML 적용: 모델 추론/재학습을 파이프라인에 직접 삽입 가능 (예: anomaly detection, 추천, LLM RAG).

관련 노트