광고 이벤트(impression · participate · reward) 데이터로 실시간 데이터 파이프라인을 구축한 기록을 위한 문서이다.
1. 왜 하는가
광고 이벤트 데이터의 실시간 파이프라인을 구축한 사례를 문서로 남기려한다.
구축 전에도 실시간 데이터를 적재하고 있었다. 말그대로 적재만 하고 있었다. 실시간 데이터를 한시간에 한번씩 bigquery 에 넣기만 하고, 그에 따른 추가적인 지표는 해당 데이터가 쌓인 다음날 일배치로 다시 생성되는 파이프라인으로 되어 있었다.
실시간으로 데이터를 보는게 아닌, 대시보드 까지 하루가 걸리는 로직이였다.
그래서 데이터를 어떻게 빠르게 분석할수 있는지에 대해서 결국 아키텍처를 뜯어 고쳐야 한다고 생각했다.
기존 파이프라인의 문제 (AS-IS)

실무 경로는 vendor postback → Kafka → Logstash → BigQuery였고, 데이터 엔지니어로써 개입 가능한 지점은 Kafka 이후였다. BigQuery는 강력하지만 실시간 지표 관점에선 이런 문제가 있었다.
| 1 | 지표가 실시간이 아님 | BQ는 배치 지향. 분 단위 캠페인 성과를 못 봄 |
| 2 | 조회 비용 ∝ 조회 빈도 | BQ on-demand scan cost. 대시보드 새로고침이 곧 비용 |
| 3 | Logstash = 테스트 불가 변환 레이어 | 복잡 변환(pub_ads explode) 부적합, 검증 실패 관측 안 됨 |
| 4 | 중복 그대로 적재 | action_id 재시도 중복 → 지표 부풀림 (계속해서 중복 적재 삭제중) |
| 5 | 퍼널 조인 비쌈 | reward가 participate보다 지연 도착(멀티스텝) → 다중 파티션 스캔 |
참고
진행 주요 데이터는 광고 이벤트(impression → participate → reward)다. 사용자가 광고를 보고(impression), 참여하고(participate), 보상을 받는(reward) 퍼널에서 PTR(참여율)·CVR(전환율)을 분 단위로 실시간 관측하는 게 이번 목표다.
2 이번 step0 에서 진행하는 작업
먼저 최종 데이터 확인을 위한 실시간 데이터 처리와 대시보드 생성이다. (앞으로 계속해서 발전 나갈 예정이다.)
개입 가능한 지점이 Kafka 이후였으므로, Kafka 뒤를 실시간 서빙에 맞게 재설계했다. 전체 흐름은 이렇다.

파일(raw 로그) → producer → Kafka(impression·participate·reward 3토픽)
→ consumer(검증·변환) → ClickHouse(실시간 rollup) → HyperDX(대시보드)
핵심 아이디어는 ClickHouse의 Materialized View다. 원본 이벤트가 들어오면 MV가 insert 즉시 분단위 집계로 접고, 대시보드는 그 집계 테이블만 조회한다. 이러면 BigQuery처럼 "조회할 때마다 원본을 스캔"하지 않으므로 — 실시간이면서 조회가 싸다.
3 진행사항
1) 인프라를 쪼개서 구성
ClickHouse / MongoDB / HyperDX로 분리로 대시보드에 필요한 관측 데이터는 ClickHouse(OLAP), 앱 상태(대시보드·유저)는 MongoDB. 용도별 저장소 분리(polyglot persistence) 시켰다.
2) producer → Kafka → consumer → ClickHouse
- producer: raw vendor 로그 파일을 Kafka 3개 토픽(impression/participate/reward)으로 리플레이
- consumer: vendor_adapter로 raw → canonical 변환. impression은 pub_ads 배열을 explode(1 노출 → N 광고 행), pydantic으로 검증, 마이크로배치로 ClickHouse 적재
- ClickHouse: ad_events(raw) + Materialized View로 분단위 rollup + PTR/CVR view
producer.py
(단순하게 impression.jsonl / participate.jsonl / rewards.jsonl 의 파일을 읽어서 카프카로 전송)
import asyncio, json, os, sys
from pathlib import Path
from aiokafka import AIOKafkaProducer
BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "localhost:9092")
DATA_DIR = Path(os.getenv("DATA_DIR", "./data/vendor"))
SOURCES = [
("impressions.jsonl", "ad.events.impression", "request_id"),
("participate.jsonl", "ad.events.participate", "request_id"),
("rewards.jsonl", "ad.events.reward", "click_key"),
]
async def main():
producer = AIOKafkaProducer(
bootstrap_servers=BOOTSTRAP,
value_serializer=lambda v: json.dumps(v, ensure_ascii=False).encode(),
key_serializer=lambda k: (str(k).encode() if k is not None else None),
)
await producer.start()
total = 0
try:
for fname, topic, key_field in SOURCES:
path = DATA_DIR / fname
if not path.exists():
print(f"[warn] missing {path}", file=sys.stderr)
continue
n = 0
with path.open(encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
row = json.loads(line)
await producer.send(topic, value=row, key=row.get(key_field))
n += 1
print(f"{topic}: sent {n}", file=sys.stderr)
total += n
finally:
await producer.flush()
await producer.stop()
print(f"done. total={total}", file=sys.stderr)
if __name__ == "__main__":
asyncio.run(main())
결과적으로 파일 → producer → Kafka → consumer → ClickHouse → HyperDX가 무손실로 흘렀다.
가장 중요한 clickhouse에서의 테이블 정보는 다음과 같다.
CREATE DATABASE IF NOT EXISTS ad;
-- 1) raw 이벤트 테이블 (canonical)
CREATE TABLE IF NOT EXISTS ad.ad_events
(
event_date Date DEFAULT toDate(event_time),
event_time DateTime64(3),
event_type LowCardinality(String), -- impression | participate | reward
event_id String, -- 합성 자연키 (mapping.md, dedup 기준)
-- 퍼널 조인 키
request_id String DEFAULT '', -- impression ↔ participate
click_key String DEFAULT '', -- participate ↔ reward
-- 광고 식별
ad_id Int64 DEFAULT 0,
campaign_id Int64 DEFAULT 0,
publisher_id Int64 DEFAULT 0,
publisher_app_id Int64 DEFAULT 0,
position Int32 DEFAULT -1, -- impression: pub_ads 배열 index (position bias용)
-- 사용자/디바이스
user_id String DEFAULT '', -- n_uid
advertising_id String DEFAULT '', -- n_advertising_id
is_lat Int8 DEFAULT -1,
platform Int8 DEFAULT 0,
device_model String DEFAULT '',
os_ver String DEFAULT '',
sdk_ver String DEFAULT '',
app_ver String DEFAULT '',
store LowCardinality(String) DEFAULT '',
-- 지역
country LowCardinality(String) DEFAULT '',
state String DEFAULT '',
city String DEFAULT '',
tab_slug String DEFAULT '',
-- participate 전용
landing_url String DEFAULT '',
-- 에러 (participate/impression의 error_message)
is_error UInt8 DEFAULT 0,
error_message String DEFAULT '',
-- reward 전용
result UInt8 DEFAULT 1, -- reward.result (true만 CVR 집계)
reward_amount Float64 DEFAULT 0, -- reward.reward (사용자 지급)
cost Float64 DEFAULT 0, -- reward.cost (광고주 과금)
reward_event_id Int64 DEFAULT 0, -- vendor reward.event_id (행 고유 아님)
event_code String DEFAULT '',
event_step Int32 DEFAULT 0,
event_set_id Int64 DEFAULT 0,
source LowCardinality(String) DEFAULT '',
publisher_app_name String DEFAULT '',
ingested_at DateTime64(3) DEFAULT now64(3), -- ReplacingMergeTree 버전 컬럼(최신 적재 우선)
INDEX idx_request_id request_id TYPE bloom_filter GRANULARITY 4,
INDEX idx_click_key click_key TYPE bloom_filter GRANULARITY 4
)
ENGINE = MergeTree
PARTITION BY toYYYYMMDD(event_date)
ORDER BY (campaign_id, event_type, event_time)
TTL event_date + INTERVAL 90 DAY;
-- 2) 분단위 rollup (SummingMergeTree) — 실시간 근사 speed layer
CREATE TABLE IF NOT EXISTS ad.ad_metrics_1m
(
minute DateTime,
campaign_id Int64,
ad_id Int64,
impressions UInt64,
participates UInt64, -- is_error=0만
participate_errors UInt64,
rewards UInt64, -- result=1 reward 이벤트 수 (멀티스텝 포함)
reward_amount Float64,
cost Float64
)
ENGINE = SummingMergeTree
ORDER BY (campaign_id, ad_id, minute);
-- 3) raw → rollup 실시간 MV
CREATE MATERIALIZED VIEW IF NOT EXISTS ad.mv_ad_metrics_1m TO ad.ad_metrics_1m AS
SELECT
toStartOfMinute(event_time) AS minute,
campaign_id,
ad_id,
countIf(event_type = 'impression') AS impressions,
countIf(event_type = 'participate' AND is_error = 0) AS participates,
countIf(event_type = 'participate' AND is_error = 1) AS participate_errors,
countIf(event_type = 'reward' AND result = 1) AS rewards,
sumIf(reward_amount, event_type = 'reward' AND result = 1) AS reward_amount,
sumIf(cost, event_type = 'reward' AND result = 1) AS cost
FROM ad.ad_events
GROUP BY minute, campaign_id, ad_id;
3) 멱등성(dedup)
producer를 실행했더니 기존 participate가 2,674 → 5,348로 정확히 2배가 되는 현상을 있었다.
물론 정상적인 상황이라면 일어나지 않을수도 있지만, 지금 현재 로직에서도 데이터를 빅쿼리에 실시간 적재시에 insert_id가 없어서 중복 적재되고 있고, 그걸 다음날 계속 해서 지워나가는 로직이 따로 돌고 있다. (결국 해당 날짜에서는 중복된 값을 볼수밖에 없는 현실이다. )
그래서 해당 중복 현상이 해결 가능한지를 테스트 해보았다.
-- 1) ReplacingMergeTree를 넣었는데도 중복이 남았다
rows_final = 2674, uniq_ids = 2213 -- 461개가 여전히 중복!
-- 원인: ORDER BY (campaign_id, event_type, event_time)
-- → dedup 키에 event_time이 있어서,
-- 같은 event_id라도 시각이 다른 "재시도 중복"은 다른 행으로 취급됨
-- 2) dedup 키를 event_id로 (event_time 제거가 핵심)
ENGINE = ReplacingMergeTree(ingested_at)
ORDER BY (event_type, event_id)
-- 결과: 5348 → 2213 (재처리 중복 + action_id 재시도 중복 모두 제거)
그런데 여기서 두 번째 함정. raw를 dedup해도 rollup(SummingMergeTree MV)은 여전히 2배였다. MV는 insert 시점에 합산하는 트리거라, 나중에 raw에서 중복이 사라져도 소급되지 않는다. raw dedup ≠ rollup dedup.
해결은 Lambda식 분리 — rollup은 빠른 근사(speed layer)로 두고, 정확한 지표는 raw에서 uniq 기반 view로 뽑았다 (uniqExact는 행 중복을 자동으로 접으므로 재발에도 면역).
4. 진행하면서 알게된 사실들
- Kafka — 선택이 아니라 주어진 조건. 기존 운영 경로가 postback→Kafka→…라 개입 가능한 지점이 "Kafka 이후"였다.
- ClickHouse — AS-IS 문제 #1·#2(실시간 아님·조회 비용)를 정면으로 푼다. Materialized View가 insert 즉시 분단위 rollup을 갱신해 분 단위 실시간 지표를 주고, 조회는 이미 접힌 rollup만 스캔하므로 BQ의 scan-cost 모델에서 벗어난다. 대안도 따져봤다 — Pinot/Druid는 컴포넌트가 많아 로컬·소규모에 과하고, TimescaleDB는 고카디널리티 집계·압축이 열세다. 현재 규모(수천 events/sec)엔 단일 바이너리 ClickHouse가 최적.
- HyperDX — ClickHouse 네이티브 연동이라 서빙 데이터를 그대로 탐색·대시보드화할 수 있다. all-in-one을 굳이 쪼갠 건 학습 목적(용도별 저장소 이해)과 자원 격리 때문.
- Python consumer — 이 규모엔 Flink가 과하다고 판단해, 팀(나) 역량이 가장 높은 Python으로 시작했다. 규모가 커지거나 exactly-once가 필요해지면 그때 Flink로 교체하면 된다.
'ML > 데이터 분석' 카테고리의 다른 글
| MMR 논문 : The Use of MMR Diversity-Based Reranking for Reordering Documents and Producing Summaries (0) | 2026.02.02 |
|---|---|
| ALS 논문 : MATRIX FACTORIZATION TECHNIQUES FOR RECOMMENDER SYSTEMS (0) | 2026.01.21 |
| [airflow] certified (0) | 2025.10.29 |
| [airflow] Branching (0) | 2025.10.27 |
| CoxPHFitter (2) | 2025.07.23 |