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 и время приёма;Ingest API аутентифицирует источник, ограничивает размер и допустимый диапазон
времени, но не объявляет событие billable. Сначала событие попадает в Kafka и в
неизменяемый архив. Решения duplicate, invalid и fraud сохраняются рядом с
версией правила, чтобы их можно было объяснить и переиграть.
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, а 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.
Checkpoint Flink защищает состояние оператора, но сам по себе не доказывает end-to-end exactly-once. Для такого результата нужны replayable source и sink, который участвует в checkpoint/transaction либо применяет идемпотентный ключ.
read_committed и применяет уникальный ключ
(advertiser_id, window, revision, kind);Это не означает, что «каждый event физически обработан ровно один раз». Повторная обработка допустима; наблюдаемый итог остаётся детерминированным.
Допустим, средняя нагрузка равна 1 000 000 попыток в секунду, а средний wire
event — 500 B:
1 000 000 × 86 400 = 86,4 млрд попыток в сутки;86,4 млрд × 500 B = 43,2 TB/сутки в десятичном измерении;15,768 PB до compression, indexes и replication;3–5× payload займёт примерно 3,15–5,26 PB;
metadata, маленькие файлы и replicas считаются отдельно;3× означает проектную проверку около 3 млн events/s, но реальный
коэффициент берётся из трафика.Память exact-dedupe нельзя оценивать одной красивой цифрой. Только raw payload ключей длиной 16 B составит:
300 млн × 16 B = 4,8 GB;3,6 млрд × 16 B = 57,6 GB;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:
Порог вроде 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 выбираются после нагрузочного теста. Ни latency, ни recall fraud model, ни compression ratio не являются гарантией самой архитектуры.
Введите числа или выберите пресет