Design a Time-Range Event Query System
Problem Design a system that continuously ingests events arriving at regular intervals and supports efficiently fetching all events falling within a given time range.
Functional requirements
- Ingest a continuous stream of timestamped events.
- Query all events within [start, end).
- Support filtering within a range by a secondary dimension (source/device/type).
- Scale storage as volume grows without degrading range queries.
- Expire or downsample old data.
Non-functional requirements
- ~500k events/sec sustained; average event ~200 bytes → ~100 MB/sec, ~8.6 TB/day raw.
- Retention 30 days hot (~260 TB, compressed ~10-20x to ~15-25 TB), 1 year cold.
- Range query p99 < 500 ms for a 1-hour window (~1.8B events — so the query must never touch most of the data).
- Write path must absorb bursts: buffer and batch rather than one row per event.
- Events may arrive out of order by up to a few minutes (network/mobile buffering).
Key components
- Ingest tier: events → Kafka (partitioned by source_id), consumed by writers that buffer and flush in batches.
- Time-partitioned storage: per-hour or per-day partitions/shards so a range query touches only the relevant partitions and old data is dropped by detaching a partition rather than by DELETE.
- Index: a B+tree on timestamp in a relational store, or an LSM-tree store with a time-prefixed key, or a purpose-built time-series DB (columnar, with delta-of-delta timestamp encoding).
- Compaction/downsampling job rolling raw events into coarser aggregates as they age.
- Cold tier: Parquet on object storage for aged data, queried on demand.
Deep dives / trade-offs
- B+tree vs LSM for this workload: a B+tree index on timestamp gives excellent range scans (leaves are ordered and linked) but suffers write amplification and page splits at 500k inserts/sec. LSM-trees absorb writes sequentially and are the natural fit, at the cost of read amplification across levels and compaction I/O. Because the timestamp is monotonically increasing, inserts land at the rightmost edge — which is either ideal (sequential, no fragmentation) or a hotspot on the newest partition, depending on how you shard.
- Partitioning scheme: partition by time and the range query prunes to a few partitions and retention is a metadata operation. But the newest partition takes 100% of writes — a single hot shard. Partitioning by (source, time) spreads writes but makes a cross-source range query a scatter-gather. This tension is the core of the design.
- Relational + timestamp index vs a purpose-built TSDB: Postgres with BRIN indexes and declarative partitioning is genuinely good here (BRIN exploits natural time-ordering for a tiny index) and keeps operational surface small. A TSDB adds columnar compression, timestamp delta encoding and built-in downsampling — a 10x storage win at 8.6 TB/day is material — at the cost of a new system and weaker ad-hoc query support.
- Batching vs latency: flushing in batches of 10k is what makes 500k/sec affordable, but it means an event isn't queryable for up to the flush interval, and an unflushed buffer is lost on a crash. Kafka as the durable buffer resolves the loss; the visibility lag remains a stated trade-off.
- Out-of-order arrival: a late event belongs in a partition that may already be compacted or downsampled. Define a watermark/grace window, and decide whether late events are dropped, backfilled, or written to a correction stream.
- Query bounds: an unbounded range ("all events") scans everything. Cap the window, require pagination, and pre-aggregate for the common dashboard queries rather than scanning raw rows.
asked …