프로젝트 소개
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 삽입/스키마 불일치 등의 일반적인 장애 모드를 다룬다.
Comments
0 Rating appears after 10 ratings
Sign in to join the discussion.