본문 바로가기

ML/데이터 분석

광고 이벤트 실시간 데이터 파이프라인 구축 — Step 0

광고 이벤트(impression · participate · reward) 데이터로 실시간 데이터 파이프라인을 구축한 기록을 위한 문서이다.

 

1. 왜 하는가

광고 이벤트 데이터의 실시간 파이프라인을 구축한 사례를 문서로 남기려한다. 

구축 전에도 실시간 데이터를 적재하고 있었다. 말그대로 적재만 하고 있었다. 실시간 데이터를 한시간에 한번씩 bigquery 에 넣기만 하고, 그에 따른 추가적인 지표는 해당 데이터가 쌓인 다음날 일배치로 다시 생성되는 파이프라인으로 되어 있었다. 

실시간으로 데이터를 보는게 아닌, 대시보드 까지 하루가 걸리는 로직이였다. 

그래서 데이터를 어떻게 빠르게 분석할수 있는지에 대해서 결국 아키텍처를 뜯어 고쳐야 한다고 생각했다.

기존 파이프라인의 문제 (AS-IS)

실무 경로는 vendor postback → Kafka → Logstash → BigQuery였고, 데이터 엔지니어로써 개입 가능한 지점은 Kafka 이후였다. BigQuery는 강력하지만 실시간 지표 관점에선 이런 문제가 있었다.

두번째로는 실시간 지표를 보기 위해서 redash-> 서비스 운영 rds에 접근하여 쿼리를 실행하고 있었다. 이러다 보니 1분마다 갱신하며 확인하는 쿼리들이 많아지면서 DB에 과부하가 걸리고 있었다.
 
문제원인

 

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. 진행하면서 알게된 사실들

 
 
Q: consumer lag이 0이 아니면 데이터 처리에 실패한 것인가?
A: 아니다. lag은 "지금 이 순간 아직 안 읽은 잔량"이지 실패가 아니다. 라이브 트래픽은 계속 들어오므로 항상 몇 개는 밀린다. 중요한 건 크기가 아니라 "시간에 따라 계속 증가하느냐"다. 늘어나기만 하면 consumer가 못 따라가는 것, 출렁이면 정상.
 
Q: unique 값으로 ReplacingMergeTree의 ORDER BY에 event_id를 "추가"하면 되나?
A: 반만 맞다. dedup은 ORDER BY 튜플 전체가 같아야 작동한다. event_time을 남긴 채 event_id를 뒤에 붙이면, 재시도(같은 id·다른 시각)는 여전히 다른 튜플이라 안 잡힌다. event_time을 키에서 빼고 event_id로 식별해야 한다 — dedup 키는 곧 "이 행의 정체성"인데, 정체성은 event_id이지 event_time이 아니기 때문.
 
Q: raw를 dedup했는데 왜 rollup은 아직 2배인가?
A: MV는 insert 트리거라 중복이 들어온 순간 이미 합산해 버렸고, raw를 나중에 정리해도 소급되지 않는다. 그래서 rollup은 speed layer(빠른 근사)로 두고, 정확값은 raw의 uniq view로 뽑는 Lambda 아키텍처로 분리했다. "속도 vs 정합성"을 한 테이블로 풀려 하지 말고 레이어를 나누는 것이 정석.
 
Q: 왜 이 서비스들(Kafka · ClickHouse · HyperDX)을 골랐나?
A: 각각 "선택"의 성격이 다르다.
  • 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로 교체하면 된다.