Design Real-Time Ad Click Event Aggregation
1. Problem Statement & Scope
System Mission
Design a mission-critical, fault-tolerant, real-time ad click event ingestion and aggregation pipeline capable of processing peak traffic of over 100,000 clicks per second (). The architecture must enforce strict Exactly-Once processing semantics for financial billing, execute multi-granularity tumbling and sliding window aggregations (1-minute, 1-hour, 24-hour CTR and spend) using Apache Flink, eliminate botnet and click-farm fraud via edge WAF and Redis sliding Bloom filters, and combine real-time stream aggregation with hourly EMR/Athena batch reconciliation (Lambda/Kappa architecture).
Synthesizing vector architecture diagram...
Functional Requirements
- Ad Click Ingestion: Ingest and validate signed click tokens globally with sub-20ms response times and HTTP 302 advertiser redirect handoff.
- Multi-Granularity Window Aggregation: Continuously compute click counts, valid clicks, bot/fraud clicks, and cumulative monetary spend partitioned by
ad_id,campaign_id, andpublisher_idacross 1-minute tumbling, 1-hour tumbling, and 24-hour sliding windows. - Strict Exactly-Once Billing Semantics: Guarantee that advertisers are billed exactly once per valid click, even across network partitions, Flink worker crashes, and upstream message replays.
- Real-Time Click Fraud Suppression: Filter out duplicate click IDs, rapid replay loops (), datacenter botnet traffic, and unverified token signatures before billing aggregation.
- Advertiser Analytics Query API: Serve advertiser dashboards with sub-50ms query responses tracking Click-Through Rate (CTR) and real-time budget burn rate.
- Lambda / Kappa Stream-Batch Reconciliation: Continuously stream raw Parquet events into an S3 data lake with hourly EMR/Athena batch jobs resolving late-arriving events and calculating billing adjustments.
Non-Functional Requirements (SLAs/SLOs)
- High Throughput: Support average , scaling to peak capacity of without data loss.
- Aggregation Freshness / Pipeline Latency: Aggregated metrics updated in advertiser dashboards within ().
- Data Correctness: financial billing accuracy with end-to-end exactly-once semantics. Zero double-billing under worker failover.
- Availability: uptime SLA for click ingestion endpoints ( downtime/year).
- Fault Tolerance: Lossless Flink recovery via RocksDB checkpoints in S3 within of node crash.
- CAP / PACELC Classification:
- Click Ingest Tier: AP / PA-EL (Availability and sub-20ms latency prioritized; immediate Kinesis buffer commit).
- Billing Aggregation & Flink State: CP / PC-EC (Strict consistency via RocksDB checkpoints, idempotent transactional DynamoDB sinks).
2. Capacity & Scale Estimation
Traffic Calculations
- Daily Ad Impressions: impressions/day.
- Average Click-Through Rate (CTR): .
- Click Ingestion Throughput:
- Query Load (Advertiser Dashboards):
- 100,000 active campaigns queried every 10 seconds by advertiser automation bots.
Storage & Data Sizing Math
- Raw Click Event Payload: (Protobuf/JSON:
click_id,ad_id,campaign_id,publisher_id,user_id,ip_hash,cpc_cents,event_timestamp,token_signature). - Columnar Snappy Parquet Data Lake (S3):
- Snappy Parquet achieves compression over raw JSON: 3-year data lake retention: .
- Aggregated Metric Record Volume (Write Reduction):
- 100,000 active ads aggregated across 1-minute tumbling windows:
- Write Load Reduction: Compared to writing raw clicks (), stream aggregation slashes database write IOPS from , achieving a write load reduction.
Memory & Cache Sizing (80/20 Pareto Working Set)
- Redis Sliding Bloom Filter (Duplicate Click Suppression):
- Sliding window of 10 minutes ().
- Number of clicks in 10 minutes at peak: .
- Target false positive probability ():
- Double-buffered rotating Bloom filter (Current + Previous window): .
- IP Velocity Rate Limiting (ElastiCache Redis):
- Top 10 Million active IPs tracked via sliding token bucket in Redis:
- Flink RocksDB Stateful Working Set:
- In-memory window state for 100,000 campaigns across 1-minute and 1-hour active windows:
Fleet Sizing & Kinesis Provisioning
- Amazon Kinesis Data Streams Provisioning:
- 1 Kinesis shard provides ingress or .
- Required shards for and :
- Apache Flink Compute Sizing (AWS Kinesis Analytics):
- Sized at 32 Kinesis Processing Units (KPUs, each with 1 vCPU and 4 GB RAM).
- Parallelism (aligned 1:1 with Kinesis shards).
3. AWS-First High-Level Architecture
Synthesizing vector architecture diagram...
Follow a click through the hot path and then the cold path. In the "Stateless High-Throughput Ingestion Fleet" panel, the NLB spreads clicks across ingest tasks, which write to Kinesis partitioned by hash(ad_id + salt), so one viral ad is spread over several shards. In the "Stateful Stream Processing" panel, Flink groups clicks by the time they happened (event time), waits for late arrivals up to a watermark, drops duplicates and suspicious IPs using Redis, checkpoints its state to S3 every 30 s, and writes per-minute totals to DynamoDB. In the "Cold Path Data Lake & Batch Reconciliation" panel, the same raw clicks are also written to S3 by hour, and an hourly Spark job recomputes the totals, correcting DynamoDB for clicks that arrived too late. Advertisers read the totals through the "Advertiser Dashboard Serving Tier" panel. The fast path gives near-real-time numbers; the batch path makes the billed numbers exact.
Data Flow Walkthrough
- Ingest & Instant Redirect: The user clicks an ad. The browser hits the NLB. The ECS Go Ingestion worker validates the HMAC impression token signature, dispatches the click event to Kinesis, and immediately issues an HTTP 302 redirect to the advertiser landing page in .
- Kinesis Partition Salting: To prevent mega-campaigns (e.g. Super Bowl ads) from overwhelming a single Kinesis shard, the partition key is salted:
ad_id + "#" + random(0, 63), spreading writes across up to 64 shards. The salt range is sized from the hottest campaign: clicks/sec on one ad divided by the records/sec shard limit needs at least 40 salts; 16 would still leave /sec per shard and throttle. - Stateful Flink Aggregation: Apache Flink consumes Kinesis events using Event-Time watermarks. It checks ElastiCache Redis Bloom filters to drop duplicate click IDs and verifies IP velocity limiters. Valid clicks increment tumbling 1-minute and 1-hour window aggregates in RocksDB.
- Idempotent DynamoDB Sink: DynamoDB is not a transactional participant in Flink's two-phase commit, so the sink must be idempotent instead: when a 1-minute window closes, Flink writes the window's absolute totals to the row keyed by
(ad_id, window_start)withSET valid_clicks = :v, spend_cents = :s, checkpoint_id = :cidand the conditionattribute_not_exists(checkpoint_id) OR checkpoint_id < :cid. A replay after a crash rewrites the same totals for the same window and is a no-op. UsingADD valid_clicks :inchere is the classic double-billing bug: every replayed flush would add the window again. - Cold-Path Reconciliation: Firehose streams raw events to S3 in Snappy Parquet. An hourly Amazon EMR job computes ground-truth aggregates over the raw data lake, catching late-arriving clicks beyond Flink's 5-minute watermark and writing financial reconciliations.
Concrete Step-by-Step Request Walkthrough: Tracing Click Ingest to Billing Update
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| 1 | User clicks ad; browser issues POST /v1/events/adclick | NLB routes TCP stream to ECS Go Ingest Task | Task validates HMAC-SHA256 click token signature and impression timestamp () | Returns HTTP 302 redirect to user in ; enqueues background record |
| 2 | Ingest task dispatches event to Amazon Kinesis | Task computes hash(ad_nike_01 + "#" + salt_7) | Writes 500-byte JSON payload to designated Kinesis shard | Kinesis assigns sequence number; returns ACK |
| 3 | Apache Flink worker reads record from Kinesis shard | Flink extracts event timestamp ; checks watermark | Watermark ; event falls within active 1-min window | Event admitted to Flink stream execution graph |
| 4 | Flink executes fraud filter check against Redis | Worker issues pipelined BF.EXISTS on click_dedupe_bf | Redis returns 0 (not duplicate); issues BF.ADD and checks IP velocity () | Click designated VALID; CPC amount added to window accumulator |
| 5 | 1-Minute Tumbling Window fires () | Flink merges RocksDB in-memory state: 840 clicks, $378.00 spend | Prepares Two-Phase Commit transaction for DynamoDB sink | Flink takes snapshot checkpoint to Amazon S3 |
| 6 | Transactional sink executes atomic DynamoDB upsert | Sink writes to AD#ad_nike_01, WIN#1M#1773648000 | Executes conditional UpdateItem writing the window's absolute totals (SET valid_clicks = :c, spend_cents = :s, checkpoint_id = :cid if checkpoint_id < :cid) | Record committed in DynamoDB in ; P99 pipeline latency |
| 7 | Advertiser dashboard queries campaign metrics | Query Service hits ElastiCache; on miss reads DynamoDB GSI | GSI CAMP#nike_spring aggregates all ad child rows for target hour | Dashboard displays live CTR () and budget burn in |
4. API Interface Design
1. Ingest Ad Click Event (POST /v1/events/adclick)
Receives incoming click beacons, validates cryptographic signatures, and outputs redirect instructions.
httpPOST /v1/events/adclick HTTP/1.1 Host: click.ads.aws.internal User-Agent: Mozilla/5.0 (iPhone; CPU iPhone OS 17_4 like Mac OS X)... X-Forwarded-For: 203.0.113.195 Content-Type: application/json { "click_id": "clk_88a91c74f0b21a", "ad_id": "ad_nike_airmax_2026", "campaign_id": "camp_nike_spring_global", "publisher_id": "pub_techcrunch_01", "user_id": "usr_998124a87", "cpc_cents": 45, "impression_timestamp": 1773648000100, "click_timestamp": 1773648002340, "click_token_sig": "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855" }
Response: 202 Accepted
json{ "status": "ACCEPTED", "redirect_url": "https://www.nike.com/air-max-2026?utm_source=ad_campaign", "ingest_timestamp": 1773648002352 }
2. Query Real-Time Campaign Performance (GET /v1/analytics/campaigns/{id}/metrics)
Returns pre-aggregated campaign performance metrics across specified time windows.
httpGET /v1/analytics/campaigns/camp_nike_spring_global/metrics?window=1h&start_time=1773648000&end_time=1773651600 HTTP/1.1 Host: dashboard.ads.aws.internal Authorization: Bearer <jwt_token>
Response: 200 OK
json{ "campaign_id": "camp_nike_spring_global", "window": "1_HOUR", "start_time": 1773648000, "end_time": 1773651600, "metrics": { "total_impressions": 5000000, "total_clicks": 142850, "valid_clicks": 139210, "fraud_clicks_blocked": 3640, "total_spend_cents": 6264450, "ctr_percentage": 2.78, "average_cpc_cents": 45.0 } }
3. Status Codes & Error Contracts
| HTTP Status | Reason Code | Error Contract Payload | Mitigation / Client Action |
|---|---|---|---|
200 OK | SUCCESS | Aggregated analytics payload | Normal completion |
202 Accepted | CLICK_ACCEPTED | Redirect URL payload | Client navigates to target URL |
400 Bad Request | INVALID_SIGNATURE | {"error": "HMAC impression signature verification failed"} | Drop event; flag potential ad tampering |
403 Forbidden | TOKEN_EXPIRED | {"error": "Click timestamp > 10m from impression"} | Discard click from billing stream |
429 Too Many Req | IP_VELOCITY_EXCEEDED | {"error": "Click rate exceeds 5/min per IP"} | Divert to fraud sink; suppress billing |
503 Service Unavail | INGEST_QUEUE_FULL | {"error": "Kinesis ingress throttling"} | Ingest task buffers locally in memory |
5. Data Models & Storage Architecture
Relational Schema (Amazon Aurora / Redshift Cold Path Analytics)
sqlCREATE TABLE ad_clicks_raw ( click_id VARCHAR(64) PRIMARY KEY, ad_id VARCHAR(64) NOT NULL, campaign_id VARCHAR(64) NOT NULL, publisher_id VARCHAR(64) NOT NULL, user_id VARCHAR(64), ip_hash CHAR(64) NOT NULL, cpc_cents INT NOT NULL, is_valid BOOLEAN NOT NULL DEFAULT TRUE, fraud_reason VARCHAR(32), event_timestamp TIMESTAMPTZ NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() ) PARTITION BY RANGE (event_timestamp); CREATE TABLE campaign_hourly_aggregates ( campaign_id VARCHAR(64) NOT NULL, window_start TIMESTAMPTZ NOT NULL, total_clicks BIGINT NOT NULL DEFAULT 0, valid_clicks BIGINT NOT NULL DEFAULT 0, fraud_clicks BIGINT NOT NULL DEFAULT 0, total_spend_cents BIGINT NOT NULL DEFAULT 0, reconciled_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), PRIMARY KEY (campaign_id, window_start) );
DynamoDB Single-Table Schema (AdAggregatesTable)
Partition Key (PK) | Sort Key (SK) | Attributes & Payloads | GSI1-PK / GSI1-SK |
|---|---|---|---|
AD#<ad_id> | WIN#1M#<epoch_minute> | valid_clicks: 840, fraud_clicks: 12, spend_cents: 37800, updated_at | CAMP#<campaign_id> / DATE#<epoch_minute> |
AD#<ad_id> | WIN#1H#<epoch_hour> | valid_clicks: 50400, fraud_clicks: 720, spend_cents: 2268000 | CAMP#<campaign_id> / DATE#<epoch_hour> |
AD#<ad_id> | WIN#1D#<YYYY-MM-DD> | valid_clicks: 1209600, fraud_clicks: 17280, spend_cents: 54432000 | CAMP#<campaign_id> / DATE#<YYYY-MM-DD> |
CAMP#<campaign_id> | BUDGET | daily_budget_cents: 10000000, current_spend_cents: 4892000, status: "ACTIVE" | — |
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~39%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.