System Design Cases
Ad Click Aggregator
Design ad click aggregator at Google Ads / Facebook Ads scale. 1M clicks/sec sustained, near-real-time aggregation (1m lag), exact billing. Pipeline: Browser/Mobile beacons -> Edge CDN -> Click Ingest API -> Kafka (5K partitions) -> Flink (Enrich -> Dedup with RocksDB -> Fraud filter with Redis -> 1m Window Aggregator) -> Druid OLAP + S3 raw + Postgres ledger. 5 scenarios: happy click flow, dedup duplicate, fraud bot filtered, hot ad partition skew, daily billing reconciliation.
Агрегатор рекламных кликов
Эта архитектура разделяет три разных результата: неизменяемый журнал входных
событий, оперативные метрики для отчётов и точный денежный ledger. Быстрый
дашборд может получить новую ревизию после late event, а счёт никогда не
исправляется незаметным UPDATE: корректировка добавляется отдельной записью.
Контракт события
Минимальный click event содержит:
- стабильный
event_id, сгенерированный до первого retry; ad_id,campaign_idили токен атрибуции;event_timeи время приёма;- версию схемы, источник и версию consent/policy;
- только необходимые для атрибуции поля, с ограниченным сроком хранения.
Ingest API аутентифицирует источник, ограничивает размер и допустимый диапазон
времени, но не объявляет событие billable. Сначала событие попадает в Kafka и в
неизменяемый архив. Решения duplicate, invalid и fraud сохраняются рядом с
версией правила, чтобы их можно было объяснить и переиграть.
Почему salting выполняется у producer
Kafka producer выбирает partition. Если сначала записать все события с ключом
ad_id, один вирусный ad уже создаст hot partition; re-key внутри consumer не
уберёт накопившийся lag. Поэтому конфигурация ingest заранее выбирает ключ
(ad_id, salt_bucket), где bucket детерминированно вычисляется из стабильного
event_id: retry попадёт в тот же shard exact-dedupe. Число bucket — измеряемая
настройка, а не магическая константа.
Точный итог требует двух стадий:
partial-aggсчитает(ad_id, salt, window).final-aggполучает все partials и складывает их по(ad_id, window).
Без второй стадии получились бы несколько несовместимых итогов для одного ad. Если salt-политика меняется, её версия входит в событие и replay-конфигурацию.
Event time, dedupe и поздние события
Окна строятся по event_time, а watermark выражает компромисс между задержкой и
полнотой. Late event внутри allowed-lateness создаёт новую детерминированную
ревизию окна. Событие за пределами этого срока не исчезает: оно остаётся в raw и
может войти в batch reconciliation и денежную adjustment.
Онлайн-dedupe хранит точные event_id на явно выбранный replay horizon. Пример
на схеме — 24 часа. Этого недостаточно, чтобы навсегда доказать уникальность;
финальный batch повторяет точный dedupe по архиву за весь расчётный период.
Bloom filter или HLL нельзя использовать как источник billable count: это
приближённые структуры. HLL уместен для метрики reach, если рядом показаны тип
оценки и её error bound.
Границы exactly-once
Checkpoint Flink защищает состояние оператора, но сам по себе не доказывает end-to-end exactly-once. Для такого результата нужны replayable source и sink, который участвует в checkpoint/transaction либо применяет идемпотентный ключ.
- partial и final публикуют ревизии в Kafka транзакционно вместе с offsets;
- Druid использует нативную Kafka ingestion и восстанавливается по offsets;
- ledger writer читает
read_committedи применяет уникальный ключ(advertiser_id, window, revision, kind); - retry ledger writer либо вставит запись один раз, либо увидит уже существующую;
- reconciliation публикует adjustment с отдельным стабильным id.
Это не означает, что «каждый event физически обработан ровно один раз». Повторная обработка допустима; наблюдаемый итог остаётся детерминированным.
Иллюстративная ёмкость
Допустим, средняя нагрузка равна 1 000 000 попыток в секунду, а средний wire
event — 500 B:
1 000 000 × 86 400 = 86,4 млрдпопыток в сутки;- raw ingress:
86,4 млрд × 500 B = 43,2 TB/суткив десятичном измерении; - за 365 дней это
15,768 PBдо compression, indexes и replication; - при фактическом compression
3–5×payload займёт примерно3,15–5,26 PB; metadata, маленькие файлы и replicas считаются отдельно; - peak
3×означает проектную проверку около3 млн events/s, но реальный коэффициент берётся из трафика.
Память exact-dedupe нельзя оценивать одной красивой цифрой. Только raw payload ключей длиной 16 B составит:
- за 5 минут:
300 млн × 16 B = 4,8 GB; - за 1 час:
3,6 млрд × 16 B = 57,6 GB; - за 24 часа:
86,4 млрд × 16 B = 1,3824 TB.
Объектные overhead, timestamps, indexes, checkpoints и replication увеличат эти числа. Поэтому retention, key encoding и количество partitions измеряются на реалистичных данных. RocksDB здесь — embedded keyed state каждого task, а не внешний RPC-сервис.
Денежная семантика
click attempt, accepted click, billable click, impression и conversion
— разные сущности. CTR нельзя применять к уже названному потоку кликов, чтобы
получить «число billable clicks». Billable policy отдельно определяет, какие
accepted events образуют начисление и по какой цене/валюте.
Streaming ledger сначала содержит provisional revisions. Независимый batch:
- читает immutable raw и версии правил;
- выполняет exact dedupe и policy evaluation;
- сравнивает результат со всеми streaming revisions;
- при расхождении создаёт traceable adjustment;
- только после этого закрывает период.
Порог вроде 0,01% может поднимать alert или требовать ручного review, но не
разрешает оставить неправильный счёт.
Отказоустойчивость и контроль
| Сбой | Безопасная реакция |
|---|---|
| Ingest ответил неясно | клиент повторяет тот же event_id; exact dedupe не начисляет дважды |
| Kafka partition недоступен | durable producer retry; API не подтверждает приём до принятой durability policy |
| Processor упал после обработки | offsets/checkpoint откатываются; выход повторяется транзакционно или идемпотентно |
| Druid task перезапущен | native Kafka indexer восстанавливает offsets; dashboard может отстать, raw не теряется |
| Ledger writer упал после commit | повтор с тем же revision key читает существующую запись |
| Fraud model деградировал | версия policy фиксируется; quarantine можно переиграть, деньги корректируются adjustment |
| Late data превысила watermark | событие остаётся в raw и учитывается независимой сверкой |
Privacy-контур задаёт purpose limitation, retention, access audit и удаление либо псевдонимизацию полей, где это допускает финансовая обязанность. Raw archive не должен превращаться в бессрочное хранилище идентификаторов.
Сценарии
Принятое событие проходит exact dedupe, fraud policy, обе стадии агрегации и порождает отдельные metric и billing revisions.
Retry с тем же event_id заканчивается disposition в архиве; downstream
aggregate и ledger не вызываются.
Подозрительное событие получает версию правила и evidence, сохраняется для review, но не становится billable.
Salt выбирается до Kafka, а final stage собирает точный итог по всем buckets.
Допустимое late event создаёт следующую ревизию и append-only delta в ledger.
Независимый batch объясняет расхождение, пишет adjustment и затем sealing mark.
Проверяемые SLO и метрики
- accepted-ingest availability и доля ambiguous responses;
- Kafka lag по partition, а не только среднее;
- event-time freshness и распределение lateness;
- duplicate/fraud/invalid rate по версии policy;
- revision count и возраст незакрытых окон;
- ledger idempotency conflicts и reconciliation deltas;
- p50/p95/p99 dashboard freshness;
- raw archive completeness, checksum и replay drill.
Значения SLO выбираются после нагрузочного теста. Ни latency, ни recall fraud model, ни compression ratio не являются гарантией самой архитектуры.
Первичные источники
- Apache Kafka, producer partitioning и consumer ownership: https://kafka.apache.org/41/design/protocol/
- Apache Kafka, transactions и граница external sink: https://kafka.apache.org/41/design/design/
- Apache Flink, end-to-end processing guarantees: https://nightlies.apache.org/flink/flink-docs-stable/docs/connectors/datastream/guarantees/
- Apache Flink, event time и watermarks: https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/
- Apache Druid, Kafka ingestion: https://druid.apache.org/docs/latest/ingestion/kafka-ingestion/
- Google Research, Mesa и versioned aggregation для критичных рекламных метрик: https://research.google/pubs/mesa-geo-replicated-near-real-time-scalable-data-warehousing/
- Apache DataSketches, approximate sketches: https://datasketches.apache.org/docs/