본문 바로가기
어플리케이션

파이썬으로 완성하는 실시간 스트리밍: Bytewax

by forward error correction Circle 2026. 7. 16.
반응형

1. Bytewax 기술이란?

왜 필요한가 — 기존 방식의 한계

스트림 처리의 사실상 표준은 Apache Flink입니다. Exactly-once, 분산 상태, 이벤트 타임 워터마크까지 모두 갖춘 성숙한 시스템이지만 본질은 JVM 엔진입니다. 파이썬 사용자가 합류하면 다음과 같은 비용이 발생합니다.

  • 이중 런타임 — PyFlink는 Beam Portability Framework 위에서 파이썬 워커를 별도 프로세스로 띄우고 gRPC로 데이터를 주고받습니다. ML 모델 inference 1건마다 직렬화·역직렬화 비용이 누적됩니다.
  • 디버깅 격리 — JVM 스택 트레이스와 파이썬 스택 트레이스가 따로 떨어집니다. NPE 옆에 KeyError가 나란히 떨어져도 둘을 잇는 콘텍스트가 없습니다.
  • 패키징 충돌 — PyTorch, transformers, polars처럼 무거운 파이썬 의존성을 컨테이너 이미지에 함께 굽는 순간 클러스터 노드의 콜드 스타트가 분 단위로 늘어납니다.
  • Faust의 종료 — Robinhood가 만든 Faust(파이썬 Kafka Streams 클론)는 사실상 유지보수가 멈췄고, faust-streaming 포크도 활성도가 낮습니다.

기술 정의

Bytewax는 파이썬 API + Rust 엔진의 분산 스트림 처리 프레임워크입니다. 내부 엔진은 Materialize·DBSP가 채택한 것과 같은 계보의 Timely Dataflow(Naiad 논문, Microsoft Research)이며, PyO3 바인딩으로 파이썬 함수를 워커 안에서 그대로 실행합니다. 별도 JVM 워커 프로세스가 없고, 사용자는 그저 pip install bytewax 한 번으로 단일 노드부터 K8s 분산 클러스터까지 같은 코드로 실행합니다.

"Python is the first-class citizen, not a guest." — Bytewax 0.18 릴리스 노트 中. PyFlink와의 본질적 설계 차이를 한 문장으로 요약합니다.

2. Bytewax 기술 특징

특징 설명
파이썬 1급 API Dataflow 객체에 map/filter/stateful_map 같은 파이썬 함수를 직접 등록. UDF 우회 없음.
Rust 엔진 Timely Dataflow 기반. 메시지 라우팅·스케줄링·상태 관리가 모두 Rust 코드.
로컬↔분산 동일 코드 단일 프로세스 실행, 멀티 워커, K8s 분산 클러스터까지 같은 dataflow 정의를 그대로 실행.
Recovery 내장 SQLite 또는 외부 스토리지에 상태 스냅샷. 장애 후 같은 명령으로 재시작 시 상태 복원.
이벤트 타임 윈도잉 Tumbling, Sliding, Session 윈도우. Watermark 기반 늦은 데이터 처리.
커넥터 생태계 Kafka, Redpanda, Kinesis, Pulsar, S3, Postgres CDC, WebSocket, File 등 공식·커뮤니티 커넥터.
ML 친화 PyTorch·HuggingFace·scikit-learn 모델을 워커 메모리에 로드하고 인라인 추론.
K8s Operator bytewax-operator로 CRD 기반 배포·롤링업데이트·재시작 관리.

3. Bytewax 동작 방식

구성 요소

컴포넌트 역할
Dataflow 사용자가 정의하는 처리 그래프. 입력·연산·출력 노드의 DAG.
Worker 하나의 OS 프로세스 = 하나의 파이썬 인터프리터 = 1 워커. GIL을 우회하기 위해 프로세스 다중화.
Timely Dataflow Runtime Rust 엔진. 워커 간 메시지 셔플, 진행상황(progress) 추적, 스케줄링.
Input Source 파티션 단위 입력. KafkaSource, S3Source, FileSource 등. 파티션이 워커에 매핑됨.
Sink 출력. KafkaSink, StdOutSink, 사용자 정의 sink.
Recovery Store 상태 스냅샷 저장소. 기본 SQLite, S3/EFS도 지원.
Operator map, filter, key_on, stateful_map, fold_window, join_named 등 연산 함수.

데이터 흐름

한 이벤트가 들어와서 결과가 나갈 때까지 일어나는 일을 단계로 풀면 다음과 같습니다.

  1. Source 파티션 분배 — Kafka partition 0~7이 있고 워커 2개라면, 워커 A=0~3, 워커 B=4~7로 균등 분배.
  2. 이벤트 수신 — 각 워커가 자신이 담당하는 파티션에서 메시지를 pull. 파이썬 객체로 디시리얼라이즈.
  3. Stateless 연산 처리 — map/filter는 워커 내부에서 그대로 처리. 셔플 없음.
  4. 키 라우팅 — key_on이 호출되면 같은 키는 같은 워커로 모이도록 Rust 엔진이 셔플. 이때 데이터는 pickle로 직렬화되어 워커 간 이동.
  5. 상태 연산 — stateful_map, fold_window 등에서 키별 상태 사전을 갱신. 상태는 워커 메모리 + 주기적 스냅샷.
  6. 윈도우 마감 — Watermark가 윈도우 종료 시점을 넘으면 결과 emit. 늦은 데이터는 정책에 따라 폐기 또는 추가 emit.
  7. Sink 출력 — 워커가 결과를 Kafka topic 등으로 produce. 백프레셔는 sink가 느려질 때 dataflow 전체에 역으로 전파.
  8. 스냅샷 — 설정된 epoch 주기로 모든 워커가 상태를 Recovery store에 기록.

4. Bytewax 구성 및 흐름도

아키텍처 다이어그램

실제 처리 흐름 (예: 클릭스트림 실시간 집계)

단계 연산 동작
1 KafkaSource topic=clicks, group_id=bw-agg에서 메시지 consume
2 op.map("parse", json.loads) bytes → dict 파이썬 객체로 파싱
3 op.filter("only_pdp", lambda e: e["page"]=="pdp") 상품 상세 페이지 이벤트만 통과
4 op.key_on("by_user", lambda e: e["user_id"]) user_id 해시로 워커 셔플 (네트워크 비용 발생 구간)
5 op.windowing.fold_window 5분 텀블링 윈도우 + EventClockConfig로 이벤트 타임 집계
6 op.map("enrich", call_model) scikit-learn 모델로 churn 확률 추론 (워커 메모리에 모델 상주)
7 KafkaSink topic=user-pdp-aggregates로 produce. 다운스트림 BI 도구가 소비

5. Bytewax 설치 방법

로컬 (개발용)

# Python 3.9 이상 필수 (3.11~3.12 권장)
python -m venv .venv
source .venv/bin/activate

pip install bytewax==0.21.*
pip install "bytewax[kafka]"           # Kafka 커넥터
pip install bytewax-redpanda           # Redpanda 커넥터 (커뮤니티)

# 버전 확인
python -c "import bytewax; print(bytewax.__version__)"

Docker

# Dockerfile
FROM python:3.12-slim

RUN apt-get update && apt-get install -y --no-install-recommends \
    build-essential librdkafka-dev && rm -rf /var/lib/apt/lists/*

WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY dataflow.py .

ENV BYTEWAX_PYTHON_FILE_PATH=/app/dataflow.py
ENV BYTEWAX_WORKERS_PER_PROCESS=2

CMD ["python", "-m", "bytewax.run", "dataflow:flow"]

Kubernetes (운영용)

# bytewax-operator 설치
helm repo add bytewax https://bytewax.github.io/helm-charts
helm repo update
helm install bytewax-operator bytewax/bytewax-operator -n bytewax-system --create-namespace

# Dataflow CRD 배포 (간소화 예시)
cat <

6. Bytewax 사용 방법

코드 — 가장 짧은 예제

# dataflow.py
import json
from datetime import datetime, timedelta, timezone
import bytewax.operators as op
import bytewax.operators.windowing as win
from bytewax.dataflow import Dataflow
from bytewax.connectors.kafka import KafkaSource, KafkaSink, KafkaSinkMessage
from bytewax.windowing import TumblingWindower, EventClockConfig

flow = Dataflow("clicks-agg")

# 1. 입력
kin = op.input(
    "kafka-in", flow,
    KafkaSource(brokers=["broker:9092"], topics=["clicks"], starting_offset="end"),
)

# 2. 파싱 + 필터
parsed = op.map("parse", kin, lambda m: json.loads(m.value))
pdp_only = op.filter("only_pdp", parsed, lambda e: e.get("page") == "pdp")

# 3. 키 라우팅
keyed = op.key_on("by_user", pdp_only, lambda e: e["user_id"])

# 4. 5분 텀블링 윈도우, 이벤트 타임 기준
clock = EventClockConfig(
    ts_getter=lambda e: datetime.fromisoformat(e["ts"]).replace(tzinfo=timezone.utc),
    wait_for_system_duration=timedelta(seconds=10),
)
windower = TumblingWindower(
    length=timedelta(minutes=5),
    align_to=datetime(2025, 1, 1, tzinfo=timezone.utc),
)
agg = win.count_window("count", keyed, clock, windower)

# 5. 출력
out = op.map(
    "serialize", agg,
    lambda kv: KafkaSinkMessage(key=kv[0].encode(), value=json.dumps(kv[1]).encode()),
)
op.output("kafka-out", out, KafkaSink(brokers=["broker:9092"], topic="user-pdp-aggregates"))

설정 — 운영 환경 변수

변수 의미 / 권장값
BYTEWAX_WORKERS_PER_PROCESS 한 프로세스 안의 워커 수. CPU 코어 수에 맞추되 GIL 영향으로 1~2가 안전.
BYTEWAX_PROCESSES 노드당 띄울 프로세스 수. K8s에서는 replicas로 대체.
BYTEWAX_RECOVERY_DIRECTORY 로컬 Recovery DB 경로. PVC 또는 EFS 마운트 권장.
BYTEWAX_SNAPSHOT_INTERVAL 스냅샷 주기(초). 10~60초. 너무 짧으면 IO 폭증, 너무 길면 복구 후 재처리 양 증가.
BYTEWAX_LOG_LEVEL debug/info/warn. 프로덕션은 info, Rust 측 로그도 같이 잡힘.

운영 시 고려사항

  • 워커 수 = Kafka 파티션 수의 약수로 맞추기. 불균형하면 특정 워커에 메시지가 몰립니다.
  • 상태 객체는 모두 pickle 가능해야 함. lambda·로컬 클래스·DB connection은 상태로 저장 불가. 모델은 직렬화하지 말고 워커 부팅 시 새로 로드.
  • Watermark wait_for_system_duration은 5~30초. 너무 짧으면 늦은 데이터 누락, 너무 길면 윈도우 emit 지연.
  • K8s에서는 Pod 종료 시 SIGTERM grace period 60초 이상. 마지막 스냅샷을 떠야 다음 부팅에서 깨끗하게 복원.
  • Prometheus 메트릭 노출 — bytewax는 OpenTelemetry 메트릭을 발행합니다. otel-collector + Prometheus 조합 권장.

7. Bytewax 자주 쓰는 명령어와 사례

명령어 용도
python -m bytewax.run dataflow:flow 단일 프로세스 실행. 개발·디버깅용.
python -m bytewax.run dataflow:flow -w 4 한 프로세스 안 워커 4개로 실행.
python -m bytewax.run dataflow:flow -p 2 -w 2 한 노드에 프로세스 2개 × 워커 2개 = 총 4 워커. 멀티프로세스 분산.
python -m bytewax.run dataflow:flow -i 0 -a host1:2101,host2:2101 수동 분산: 노드 인덱스 0, 클러스터 주소 명시. K8s 없을 때 활용.
python -m bytewax.run dataflow:flow --recovery-directory=./recovery 로컬 Recovery 경로 지정. 같은 경로로 재실행하면 상태 복원.
kubectl get dataflows -A bytewax-operator로 배포된 Dataflow CRD 목록 조회.
kubectl logs -l app=bytewax-clicks-agg -f 워커 Pod 로그 실시간 추적. epoch 진행 상황 모니터링.

현장 사례

  • 이상거래 탐지(FDS) — Kafka로 들어오는 결제 이벤트를 사용자별 키로 윈도우 집계 → XGBoost 모델 인라인 추론 → 의심 거래만 다운스트림 Kafka로 emit. JVM/파이썬 브리지 없이 P99 200ms 이내 처리.
  • 실시간 임베딩 — 사용자 클릭 이벤트를 받아 HuggingFace sentence-transformer로 임베딩 생성 → Qdrant/Milvus로 upsert. 모델은 워커 부팅 시 1회 로드, 추론은 인라인.
  • 로그 라우팅 — Vector·Fluent Bit에서 들어온 JSON 로그를 파이썬 정규식으로 분류·익명화 → S3 + Quickwit 분기 적재. Logstash GROK 패턴을 파이썬 코드로 대체.
  • CDC 후처리 — Debezium이 만든 변경 이벤트를 받아 키별 최신 스냅샷 dict로 유지 → Postgres → 검색 인덱스 동기화.

8. Bytewax 활용 방안 — 비교와 트레이드오프

대안 기술 비교

기술 언어 지향 강점 약점
Bytewax Python 1급 파이썬 ML/데이터 생태계 직결, 단일 코드로 로컬↔분산 초고처리량(M/s) 영역에서는 Flink/Spark 대비 한계
Apache Flink JVM (Java/Scala) 성숙, exactly-once, 초고처리량, 풍부한 SQL PyFlink는 별도 워커 + 직렬화 오버헤드, 운영 복잡도
Spark Structured Streaming JVM (PySpark 가능) 배치·스트리밍 통합, 거대한 생태계 마이크로배치 모델로 초저지연 불가, 자원 무거움
Kafka Streams JVM 전용 Kafka 클러스터에 라이브러리로 임베드, 운영 단순 Kafka 토픽이 유일한 입출력, 파이썬 사용 불가
Faust / faust-streaming Python 파이썬 네이티브 (개념 선구자) 사실상 유지보수 종료, 분산 모델 약함
Quix Streams Python Kafka 중심, Pandas DataFrame 친화 API Kafka 외 입력 제한적, 분산 처리는 Kafka 컨슈머 그룹에 의존
Arroyo SQL + Rust SQL 기반, 초저지연, Rust 엔진 파이썬 UDF/ML 친화도 낮음, 생태계 초기

언제 쓰면 안 되는지

  • P99 1ms 이하 초저지연이 SLA — 파이썬 GIL 영향이 누적됩니다. 거래소 매칭엔진은 C++/Rust.
  • 이미 성숙한 Flink 운영 조직 — 굳이 갈아탈 필요 없음. 새 파이프라인부터 Bytewax 평가.
  • SQL이 1차 인터페이스인 팀 — RisingWave, Materialize, Arroyo가 더 적합. Bytewax는 파이썬 코드가 1차 표현 수단.
  • 초당 메시지 1M 이상 단일 토픽 — 가능은 하지만 Flink 대비 클러스터 자원 더 듭니다. 비용 비교 권장.
  • 완벽한 Exactly-Once 트랜잭션이 모든 sink에 보장돼야 함 — 현재 Bytewax의 EOS는 Kafka 트랜잭션 sink 등 제한된 경로에서 동작.

트러블슈팅

증상 원인 해결
워커 1개만 CPU 100%, 나머지는 놀고 있음 key_on 키 분포가 편향. 예: user_id=0이 50% 차지 키 앞에 salt 붙이기 (user_id + ":" + str(hash%N)). 집계 단계 2단으로 분리.
PicklingError: Can't pickle local object stateful_map의 builder 안에서 lambda나 로컬 클래스 반환 최상위 함수/클래스로 옮기기. 또는 dict/dataclass로 상태 구성.
Recovery DB 디스크 폭증 상태 키가 무한 증식 (TTL 없음) stateful_map에서 None 반환해 상태 삭제, 또는 session window로 자연스럽게 만료.
윈도우가 영원히 emit 안 됨 EventClockConfig의 ts가 미래거나 wait_for_system_duration이 너무 길다 소스 클럭 점검, wait_for_system_duration 5~30초로 축소, fallback에 SystemClockConfig.
Kafka lag 누적 + 워커 OOM sink가 느려서 백프레셔, 메모리에 메시지 적체 sink batch 크기·linger.ms 조정, 워커당 메모리 limit 상향, max.poll.records 축소.
SIGTERM 직후 재기동 시 같은 메시지 재처리 grace period 안에 스냅샷 못 떴음 Pod terminationGracePeriodSeconds 90~120초, 스냅샷 주기 30초로.
HuggingFace 모델 로딩에 5분 워커마다 모델을 디스크/HF Hub에서 매번 받음 이미지에 모델 가중치 포함하거나, PVC에 캐시 마운트. transformers cache_dir 통일.
Operator 배포 후 Pod 무한 재시작 Recovery 스토리지(S3) IAM 권한 미부여 IRSA 또는 ServiceAccount에 s3:GetObject/PutObject/ListBucket 부여.
로컬에선 잘 되는데 분산에서 KeyError 전역 변수에 상태 보관 (워커마다 별도) 반드시 stateful_map / fold_window로 키 단위 상태 관리. 모듈 전역은 read-only로만 사용.

주의해야 할 점

  • "파이썬이라 느릴 것"이라는 선입견을 검증하라. stateless map/filter의 CPU 사이클은 거의 Rust에서 돈다. 진짜 병목은 JSON 파싱·ML 추론 같은 사용자 코드.
  • 워커는 프로세스다, 스레드가 아니다. 모듈 임포트 비용, 모델 로딩 비용이 워커 수만큼 곱해진다. 컨테이너 이미지를 가볍게 유지.
  • 키 카디널리티를 항상 의식하라. user_id 같은 키는 무한 증식하기 쉽다. TTL 없는 stateful_map은 시한폭탄.
  • watermark는 데이터에 거짓말을 하지 않는다. 늦은 데이터가 많다면 윈도우 emit이 늦어지는 게 정상. wait 시간을 줄이면 정확도가 떨어진다.
  • 스냅샷 주기는 IO와 복구비용의 트레이드오프. 30~60초가 보통의 균형점.
  • 로컬 테스트는 -w 1로 시작하고, 그 다음 -w 4로 분산 시나리오를 강제로 만든다. 로컬 단일 워커에서 우연히 통과하던 코드가 분산에서 깨지는 일이 자주 있다.
  • 버전을 고정하라. Bytewax는 0.x 시리즈에서 API 변경이 있었다. requirements.txt에 bytewax==0.21.* 형태로 마이너 잠금.
  • 오피저/매니지드 vs DIY 선택. 운영 인력이 부족하면 Bytewax Operator 또는 매니지드 옵션. K8s 운영 역량 있으면 헬름 직접 관리.
반응형