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.
Спроектировать ad click aggregator уровня Google Ads / Facebook Ads: принимать миллионы кликов в секунду, агрегировать per-ad-per-minute / hour / day, предотвращать фрод (дедуп, бот-фильтр), выставлять рекламодателю точный счёт. Главная боль — биллинг должен быть exactly-once: пере-зачёт = жалобы регулятору, недо-зачёт = потеря выручки.
Functional:
(ad_id, user_hash, ts, geo, device, ip, referrer, cost_micros).Non-functional:
| Метрика | Значение | Расчёт |
|---|---|---|
| Clicks/day | ~100B (всех событий) | ~1B billable clicks (CTR 1%) |
| RPS avg | 1.2M events/s | 100B / 86400 |
| RPS peak | 5M events/s | 4× spike |
| Avg event size | 500 bytes | (ad_id, user_hash, ts, geo, ...) |
| Daily ingest | ~50 TB | - |
| Yearly raw storage | ~18 PB | columnar 3-5× compression |
| Kafka cluster | 200+ brokers, 5K parts | partition by ad_id |
| Flink slots | 1000+ task slots | - |
| Druid cluster | 500+ historical nodes | real-time + historical tier |
| Aggregation window | 1m (tumbling) → 1h → 1d | - |
| Dashboard QPS | 10K | - |
| Reports p95 | < 300ms | OLAP query |
| Fraud check latency | < 100ms inline | online ML signals |
Browser/Mobile -> Edge CDN -> L7 LB -> Click Ingest API -> Kafka (clicks topic, 5K partitions)
|
v
Flink: Enrich -> Dedup -> Fraud -> Window Agg (1m)
| | | |
+-> S3 raw | RocksDB Redis +-> Druid (OLAP)
| (audit) | event_id features +-> Postgres ledger
| | (billing)
+------ Checkpoint S3 -----------+
|
Reports API --> Druid + Ledger + S3 (reconcile)
Happy path: клик летит browser → edge → ingest API → Kafka. Flink consume → enrich (добавляет campaign_id, advertiser_id) → dedup (RocksDB lookup по event_id, miss) → fraud filter (Redis feature lookup, signals OK) → 1-min tumbling window aggregator. По истечении окна emit segment в Druid + INSERT в billing ledger. End-to-end < 60s.
Dedup duplicate: beacon retry после network blip отправляет тот же event_id дважды. Второй заход доходит до dedup, который смотрит RocksDB — event_id уже видели 30s назад. Дроп, но в S3 raw всё равно пишем (для аудита). Counter в Druid не инкрементится, ledger не трогается — без двойной оплаты.
Fraud bot: click farm шлёт 10K кликов/сек с одного IP, headless Chromium UA, без impression context. Fraud filter (online ML model + Redis feature store) флагает: rate > 5K/s + datacenter ASN + UA pattern. Дроп. Событие тэгается fraud=true и идёт в S3 для analytics (но не в billable count). Advertiser CTR не раздут, $0 charged.
Hot ad partition skew: виральная реклама SUPER_BOWL_2026 получает 500K rps на одну партицию (100× средней). Kafka партиция #1247 backlog растёт до 3M сообщений, p99 latency до 8s. Mitigation: producer переключается на sub-key salting (ad_id + bucket(0..31)), нагрузка размазывается по 32 sub-партициям. Aggregator scales out 10× реплик. Lag падает обратно до <1m.
Daily reconciliation: в t+24h Spark batch job читает S3 Iceberg raw для прошлого дня, пере-агрегирует независимо от streaming, сравнивает результат с billing ledger. Обнаружен delta $1,247 на advertiser_id=acme (0.003% от их spend). В пределах SLA tolerance (<0.01%) — emit adjustment txn в ledger, ledger sealed для billing run. Над тхreshold — алерт + manual review.
Контекст: Нужно считать clicks (точно — это деньги) и unique_users (приблизительно — для CTR/reach метрик в дашборде).
Опции:
| Подход | Точность | Latency | Storage | Когда |
|---|---|---|---|---|
| Flink exact aggregation | 100% | 1m | O(N events) | Billing, click counts |
| HLL/Theta sketch | ~1% err | <100ms | O(log N) | Dashboards, unique_users |
| Lambda (batch + stream) | 100% eventually | 1m + 24h | dual storage | Когда нужны оба + reconcile |
Решение: Hybrid. Flink exact aggregation для billable counts (clicks, spend) — это деньги, ошибаться нельзя. Theta sketch в Druid для unique_users / reach — погрешность 1-2% невидна в дашборде, но storage и latency на 2 порядка лучше. Daily reconciliation против S3 батчем как safety net (ADR-001 = lambda-light: stream owns truth, batch verifies).
Последствия: Сложность two paths, но сохранение exact billing + дешёвые approximate метрики. Альтернатива (всё approximate) проигрывает регулятору; всё exact — слишком дорого по storage и query latency.
Контекст: Дедупликация по event_id нужна, потому что beacons ретраятся при сетевых ошибках, мобильные SDK буферизуют события offline и могут флашить дубли.
Опции:
| Window | RocksDB state size | False negatives (миссы) | Late event handling |
|---|---|---|---|
| 5 минут | ~1B events × 16 bytes = 16 GB | High при offline mobile (часовые буферы) | Side-output для late |
| 1 час | ~12 × 16 GB = 192 GB | Medium — покрывает большинство retries | Late tolerance OK |
| 24 часа | ~4.6 TB state | Low — полностью покрывает | Самый дорогой checkpoint |
Решение: 5-минутное окно для inline dedup в Flink (16 GB state влезает в RocksDB на воркере, checkpoint быстрый). + secondary daily dedup в batch reconcile для случаев offline-buffered mobile (event приходит через 6 часов после клика). Late события (старше 5m) идут в side-output Kafka topic → batch dedup при ежедневной reconciliation.
Последствия: В streaming path возможны редкие дубли от late mobile events (<0.05%), но они ловятся в daily batch и corrected в ledger до billing run. Storage в 12-30× дешевле, чем full 1h/24h state. Альтернатива (24h) даёт 100% inline accuracy, но требует TB+ state на воркер и убивает checkpoint latency.
| Failure | Mitigation |
|---|---|
| Kafka rebalance double-count | Flink end-to-end exactly-once: checkpoint barrier + transactional Kafka producer + idempotent Druid sink |
| Flink worker OOM mid-window | Restore from last checkpoint (RocksDB на S3), replay Kafka offsets, no double-count |
| Late event past watermark | Allowed lateness 1h, side-output для еще более поздних, daily batch reconciliation |
| Fraud false positive | Soft-flag вместо hard-drop, separate "verified clean" stream для billing, manual review queue |
| Hot ad partition skew | Sub-key salting в producer (ad_id + bucket), pre-aggregate в client SDK (1s batch) |
| Druid segment merge OOM | Smaller segment granularity (1h vs 1d), separate real-time vs historical tier |
| Cost discrepancy advertiser-side | Daily reconciliation против S3 raw, immutable audit log, SLA tolerance documents |