상위: Papers
원문 링크
- doi: (미등록)
- arxiv: (미등록)

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).
관련 노트
- 논문: KARE, LightRAG, MedGraphRAG, ODKE+
- 개념: Streaming Architecture, Graph RAG, Knowledge Graph, RAGOps
- See also: ApeRAG, N-RAG Manual