Flink Window Aggregator Lab (Interactive)
Collapse 500k clicks/sec into one window row (99.9% write cut) and price lateness. Aggregate raw ad clicks into tumbling-window summary rows in Flink and quantify event-time lateness revenue loss.
Ad-Click Aggregation: Flink Tumbling Windows vs Raw Writes
Collapse 500.0k clicks/sec into one summary row per window — and see what event-time lateness costs the bill.
Writing one row per click would push 500.0k INSERTs/sec into the OLAP store; a 60s tumbling window collapses the 30.00M clicks in each window into 50,000 summary rows — a 99.83% write reduction while keeping dashboards sub-second, because ClickHouse only reads the clicks and total_cost columns (columnar + SIMD, < 50 ms over 1B rows). The billing correctness lives in event-time semantics: subway-buffered clicks arrive minutes late, and without a watermark they fall into the wrong minute and get charged (or dropped) — here that is 0 clicks and $0 a day at stake. Exactly-once rides on Flink checkpointing to S3 plus ReplacingMergeTree re-merging duplicate window rows keyed by (ad_id, campaign_id, window_minute).
How It Works Under the Hood
Writing every ad click as a database row would push hundreds of thousands of inserts per second into the OLAP store, but dashboards do not need that — they need sums per ad, campaign, country, and device per minute. Flink tumbling windows accumulate the raw stream and emit exactly one summary row per group per window, a roughly 99.9% write reduction, and ClickHouse columnar ReplacingMergeTree answers billion-row aggregates in under 50 ms. Billing correctness rides on event time: subway-buffered clicks arrive minutes late, and without a watermark holding the window open they land in the wrong bucket and get mis-counted — real dollars lost or double-charged.
Core Architectural Principles
- Write reduction = 1 - (agg_keys / (clicks/sec x window_sec)), about 99.9% at 500k/sec over 60s windows.
- An event-time watermark keeps a window open, for example plus 10 minutes, so late clicks amend their correct bucket.
- Exactly-once: Flink checkpointing plus ClickHouse ReplacingMergeTree keyed by ad, campaign, and window.
Lead with do not store raw clicks — aggregate in the stream, quantifying the roughly 99.9% write reduction that makes ClickHouse viable. Then own event-time semantics: subway and offline late clicks and the watermark versus allowed-lateness tradeoff, because mis-billing is fraud and revenue, not just inaccuracy. Mention exactly-once via checkpointing plus an idempotent sink, and a late-event side-output for delta adjustments.
Tumbling-window aggregation cuts write load by about 99.9% but adds windowing state and late-data handling that can mis-bill if event-time semantics are skipped.