프로젝트 소개

RealTime-ECommerce-Analytics(RTEA)는 Python으로 작성된 고성능 스트리밍 파이프라인으로, REST API에서 이커머스 클릭스트림 데이터를 수집해 Confluent Cloud Kafka를 통해 ClickHouse에 라우팅하고 실시간 분석 쿼리를 제공한다. ## 아키텍처 파이프라인은 다음 네 가지 레이어로 구성된다: **데이터 소스** — Railway.app에서 호스팅되는 시뮬레이션 클릭스트림 생성 API로, 초당 약 10,000개의 JSON 메시지 생성. **메시지 브로커** — Confluent Cloud Kafka, `clickstream-events` 토픽(6개 파티션), SASL_SSL 인증, Snappy 압축, 7일 유지. **처리 계층** — `kafka-python` 기반 프로듀서 서비스(16 KB 배치, gzip 압축, 32 MB 버퍼, 초당 10K 메시지로 제한) 및 컨슈머 서비스(1,000개 이벤트 배치 처리, 오프셋 자동 커밋, 필드 변환). **분석 스토리지** — ClickHouse Cloud, MergeTree 엔진 테이블 `(timestamp, user_id)` 순서, LZ4 압축, 30일 TTL. ## 데이터 흐름 1. 프로듀서는 1초마다 외부 API를 폴링하여 메시지를 큐에 수집하고, Kafka 토픽으로 배치 게시한다. 2. 여러 컨슈머 인스턴스가 할당된 파티션을 구독하고, JSON 직렬 해제, timestamp 및 필드 변환 후 ClickHouse에 배치 기록한다. 3. 사용자는 SQL을 통해 `ecommerce.clickstream_events` 테이블을 조회한다. ## 주요 설정 ### 프로듀서 - 배치 크기: 16 KB - 압축: gzip - 버퍼 메모리: 32 MB - 재시도: 3회 ### 컨슈머 - 배치 크기: 1,000개 이벤트 - 배치 타임아웃: 5초 - 오프셋 자동 커밋 활성화 ### ClickHouse 테이블 ```sql CREATE TABLE clickstream_events ( user_id String, session_id String, timestamp UInt64, event_type String, product_id String, product_category String, price Float64, quantity Int32, source String, received_time DateTime DEFAULT now() ) ENGINE = MergeTree() ORDER BY (timestamp, user_id) ``` ## 성능 결과 - 지속 처리량: 초당 650+ 이벤트 - 총 처리 이벤트: 510,414+ - 서브초 엔드투엔드 지연시간 ## 빠른 시작 1. `uv sync`로 의존성 설치 2. `.env.example`을 `.env`로 복사하고 Confluent Cloud 및 ClickHouse 자격증입력 3. 프로듀서 실행: `uv run python services/producer/src/simple_producer.py` 4. 컨슈머 실행: `uv run python services/consumer/src/simple_consumer.py` 5. 데이터 확인: `uv run python check_clickhouse.py` ## 문제 해결 저장소에는 진단 스크립트(`debug_kafka.py`, `test_clickhouse_insert.py`)가 포함되어 있으며, API 연결 시간 초과, Kafka 컨슈머 지연, 메시지 직렬 해제 오류, ClickHouse 삽입/스키마 불일치 등의 일반적인 장애 모드를 다룬다.