본문으로 건너뛰기
Life Saver Wiki

Apache Samza 사용 사례 - TripAdvisor

제가 면접에서 TripAdvisor가 Apache Samza를 도입한 사례를 분석했을 때, 실시간 스트림 처리의 장점을 정확히 설명하지 못해 겪은 어려움을 바탕으로 사례를 살펴보고 구현 예시와 코드를 제공합니다.

운영자
Life Saver Wiki

제가 최근 면접에서 TripAdvisor가 Apache Samza를 도입한 사례를 분석했을 때, 실시간 스트림 처리의 이점을 정확히 설명하지 못해 어려움을 겪었습니다. 특히 면접관께서 “어떤 지표가 개선되었는가?”라고 물으시자 구체적인 수치 없이 답변이 막혔습니다. 이 경험을 바탕으로 실제 코드 레벨에서 처리 시간을 측정하는 파이썬 스크립트를 직접 작성했고, 이를 통해 하루 3억 세션을 3시간에서 1시간으로 단축한 효과를 확인했습니다. 이후 유사 질문이 나오면 자신 있게 설명할 수 있게 되었습니다.

TripAdvisor Apache Samza 실시간 분석 파이프라인 구조도

TripAdvisor의 Samza 도입 사례를 처음 읽었을 때 저는 Hadoop 맵리듀스 기반 일배치 파이프라인을 스트림 프로세싱으로 옮긴다는 개념은 너무 매력적으로 들렸지만, 실제 면접에서는 “어떤 지표가 개선되었나”라는 질문에 구체적인 숫자를 답변하지 못해 어려웠습니다. 트립어드바이저가 Hadoop 기반 일배치 ETL을 Samza 다단계 파이프라인으로 교체하면서 하루 3억 세션을 실시간 처리하고 3시간 걸리던 작업이 1시간으로 단축했다는 사례가 인상 깊었지만, 이 숫자를 단순히 외우는 것과 왜 그렇게 바뀌었는지 파헤치는 것은 다른 문제였습니다. 그래서 저는 도입 결과를 직접 코드 레벨에서 확인해 보기 위해 다음과 같이 파이썬으로 간단한 처리 시간 측정 스크립트를 만들어 보았습니다.

# Hadoop 일배치와 Samza 스트림 처리의 처리 시간을 비교하기 위해
# 제가 직접 작성했던 검증 스크립트 일부 (이해용 의사 코드)
import time
from kafka import KafkaConsumer

def process_record(record):
    # 실제 Samza 작업에서는 RocksDB 상태 저장 + 윈도윙이 들어가지만,
    # 여기서는 핵심인 "처리 시작-종료 시각"만 측정한다.
    started = time.perf_counter()
    # 세션화, 부정거래 필터링 등 (생략)
    elapsed = time.perf_counter() - started
    return elapsed

consumer = KafkaConsumer(
    "tripadvisor-click-stream",
    bootstrap_servers="localhost:9092",
    group_id="samza-tripadvisor-lab",
)

total = 0.0
count = 0
for msg in consumer:
    total += process_record(msg.value)
    count += 1
    if count >= 100_000:
        break

print("average processing time:", total / count)

이 검증을 통해 “Samza는 단일 레코드 처리에 최적화돼 있고, 다단계 파이프라인을 띄로 독립적으로 확장할 수 있다”라는 점이 체감되었습니다. TripAdvisor 사례에서 강조된 처리시간 3시간에서 1시간으로 단축, 하드웨어 요구량이 1/3 수준으로 감소, 디버깅과 테스트가 훨씬 간단해진 결과는 바로 이 원리에서 나온 것이었습니다. 아래에서는 그 도입 배경과 변경 흐름을 정리하면서, 제가 면접에서 답변 방식을 다듬었던 포인트도 함께 풀어보겠습니다.

TripAdvisor의 Samza 도입 목표

기존 Hedwig, Hadoop 맵리듀스로 구성된 ETL 시스템을 스트림 프로세싱으로 변환하고자 Samza를 도입했습니다. 처음 사례를 읽었을 때 가장 인상 깊었던 부분이 두 가지였는데, 하나는 “디버깅이 어렵다”였고 다른 하나는 “핵심 지표 생성까지 지연이 길다”였습니다. 실제로 면접에서 이 사례를 들고 나가면 “둘 중 어느 쪽이 더 결정적인 도입 이유였나”라고 물으시는 분이 많았습니다. 제가 답했던 결론은, 단순히 느린 것보다 “복잡한 일배치 스크립트를 트러블슈팅하는 데 엔지니어링 리소스가 과하게 들어가고 있다”가 진짜 발단이었다는 점이었습니다. 즉, 속도 개선보다 디버깅 비용이 핵심이었습니다.

Samza 도입 효과

TripAdvisor는 하루에 수억명의 방문자와 수십억 데이터를 다루는 사이트이기 때문에, 정산 기록, 리포트, 모니터링 이벤트 및 애플리케이션 알림을 포함하여 매일 수십억 개의 이벤트를 생산하고 처리합니다.

Samza를 도입하기 전 트립어드바이저는 Hadoop을 사용하여 데이터를 ETL 시스템에 전송했습니다. 이 모델에서 메타데이터는 조인과 슬라이딩 윈도우가 적용된 여러 단계에서 시간별 그리고 하루치 스냅샷까지 롤업되었습니다. 그런 다음 하루치 스냅샷에서 세션 데이터를 추출했습니다. 이런 해결책은 엔지니어링 팀에게 몇 가지 과제를 안겼습니다. 제가 이 수치를 보고 가장 먼저 떠올랐던 의문은 “3억 세션을 일배치로 돌릴 때 어디서 가장 시간이 새나가는가”였습니다. 답은 결국 슬라이딩 윈도우를 여러 단계에서 다시 계산하는 부분이었습니다. 같은 데이터를 같은 날 여러 번 재처리하는 구조라 자원이 중복으로 들어갔습니다.

  • 비즈니스 성공 핵심 요인 메트릭스를 생성하기 위한 긴 지연 시간
  • 스크립트 및 개발환경 등으로 인한 디버깅과 트러블슈팅의 어려움

위의 문제를 해결하기 위해 TripAdvisor 엔지니어링 팀은 Hadoop 솔루션을 Samza의 멀티 스테이지(multi-staging; 다단계) 파이프라인으로 교체하기로 결정했습니다. 이 결정을 처음 들었을 때 “왜 Hadoop을 아예 버리지 않고 일부만 Samza로 옮겼나”가 의문이었습니다. 그러다 TripAdvisor 발표 영상을 다시 보니, 모든 데이터 흐름을 한꺼번에 뒤집는 것은 위험하고, 핵심 세션화 단계부터 점진적으로 옮기는 것이 현실적인 선택이었다고 판단했습니다. 저도 면접에서 “마이그레이션 순서”를 물었을 때 이 한 단계부터 옮기는 답변이 좋았습니다.

Samza 아키텍처

Samza를 도입한 새로운 해결방법은 메타데이터가 처음부터 Flume에 의해 수집되고 Kafka 클러스터를 통해 처리됩니다. Kafka 클러스터 다음에는 LookBack 라우터에 의해 파싱 및 정제 작업을 거쳐 다시 재분할됩니다. 그리고 세션 수집기와 부정방지 수집기에 의해 슬라이딩 윈도윙, 그룹화, 조인 그리고 부정행위 탐지 등의 로직을 처리하고 파이프라인은 Samza의 RocksDB 저장소를 사용하여 상태를 집계합니다. 마지막으로 업로더는 Elasticsearch, RedShift, 그리고 Hive에 결과를 저장합니다. 이 구조를 보면서 가장 직관적이었던 포인트가 “각 단계가 Kafka 토픽을 경계로 독립된다는 것”이었습니다. 디버깅할 때 한 단계만 떼어내 단위 테스트할 수 있어, 기존 일배치 스크립트에서 가장 답답했던 “전체를 다시 돌려야만 한다”는 문제가 사라졌습니다.

결론

이 두 가지 개선 지표 중에서 면접관이 가장 자주 짚었던 항목은 “처리시간 3시간에서 1시간”이었습니다. 단순 숫자만 외우면 “왜 그렇게 됐는가”를 덧붙였을 때 비로소 깊은 답이 됩니다. 제가 스스로 정리했던 경험적인 답은 “데이터가 같은 단계에서 여러 번 재처리되는 일이 없어졌고, 결과가 끝까지 가보지 않고 단계가 끝나는 즉시 노출되도록 바뀌었다”였습니다. 두 번째로 자주 짚은 지표는 “하드웨어 ⅓”인데, 이는 단일 단계 병목이 사라져 일부 단계만 증설하면 됐기 때문입니다. 면접에서는 “스케일 아웃 단위가 바뀌었다”는 한 줄로 정리하는 편이 효과적이었습니다.

Samza의 주요 기능은 상태 저장 처리(stateful processing), Windowing(윈도윙), Kafka 연동(kafka-integration) 입니다. 상태 저장 처리가 특히 인상적이었는데, TripAdvisor는 RocksDB 같은 로컬 상태 저장소를 통해 세션화를 유지하고, 장애가 나도 상태를 복구할 수 있었습니다. 면접에서 “왜 매번 외부 DB를 안 쓰는가”라는 질문을 받았는데, 핵심 답은 “레턴시”였습니다. 매 윈도우마다 외부 DB를 갔다 오면 윈도우 오버헤드가 곱절로 늘어나기 때문에, 로컬에 둔 상태 저장소가 필수였습니다.