0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Time-Series and Metrics Stores
Metrics look like they should fit anywhere. A timestamp, a name, a few labels, a number. Teams start with a table, and it works until it does not, usually around the point where the table holds a few billion rows and a dashboard covering the last six hours takes ninety seconds.
The reason is that the workload has a shape unlike anything a general-purpose database is tuned for, and every one of its properties points the same direction.
Writes are relentless, appends, and almost always in time order. The ingest rate is fixed by the fleet rather than by user behavior:
For machines reporting metrics at , that is 133,333 samples a second, continuously, forever, with no diurnal relief. There are no updates and there are effectively no deletes, only whole windows expiring.
Reads are the opposite shape. A query asks for a narrow time range across many series and then aggregates: "p99 latency by endpoint over the last hour." It almost never asks for one sample by primary key, which is exactly the access pattern a B-tree is built to serve.
Recent data is what anyone looks at. Essentially all queries touch the last few hours. Data older than a month is read by a compliance report or a capacity planner, rarely and in bulk.
Every value is small and the metadata around it is large. The sample is eight bytes. The series identity, meaning the metric name and its labels, is hundreds of bytes and is the same for every sample in the series.
Why the obvious relational schema fails
A table of (timestamp, metric, labels, value) with an index on time fails for three compounding reasons.
Per-row overhead dominates. In PostgreSQL a row carries roughly 24 bytes of header before any column data, so storing an 8-byte float costs several times that, and the label set is repeated in every row unless it is normalized into a join that every query then has to pay for.
Index maintenance is the write bottleneck. Every insert updates a B-tree, and a B-tree under a continuous high-rate append suffers page splits and write amplification. The write-optimized structure for this is an LSM tree, whose behavior is covered in storage engines and data structures, and that is what most time-series engines use underneath.
Range scans read the wrong bytes. A dashboard wants one column, the value, across a million rows. A row store reads every column of every row to get it.
What a time-series engine does instead
Store each series as a column of values, chunked by time. Samples for one series over a two-hour window become one compressed block. A query for the last hour opens a handful of blocks and reads nothing else, and because the values in a block are similar to one another, the compression is extraordinary, which is the subject of the next section.
Separate the series index from the sample data. Labels are stored once per series in an inverted index mapping label pairs to series ids, exactly the structure described in the search lesson. Resolving {service="checkout", status="500"} is a postings intersection producing a set of series ids, and only then are those series' chunks read.
Partition by time, and make expiry a drop. Chunks live in time-bounded partitions, so retention is implemented by deleting whole partitions rather than by a delete of rows, which would be the single most expensive operation in the system. The general argument for time-based partitioning is in partitioning and sharding; time-series stores are its purest case, since the partition key is the query filter for essentially every query.
The data models, and the vocabulary
Three models are in wide use and the differences matter more than they look.
The labeled model, used by Prometheus and the ecosystem around it, identifies a series by a metric name plus a set of key-value labels, and every distinct label combination is a distinct series. It is simple and it puts all the pressure on cardinality, which is the subject of a later section.
The measurement model, used by InfluxDB, splits the identity into indexed tags and unindexed fields, so several related values can share one timestamp and one identity. The tag and field distinction is a real design decision that people get wrong, since putting a high-cardinality value in a tag is the same fatal mistake as making it a label.
The relational model, used by TimescaleDB, keeps SQL and a hypertable that is transparently partitioned by time. You get joins, real transactions, and every tool that speaks Postgres, at the cost of a less specialized engine.
The metric types are worth stating precisely because misclassification produces wrong dashboards.
A counter only increases and resets to zero on restart. You never graph its value, you graph its rate, and the query engine has to detect resets.
A gauge goes up and down and its instantaneous value is meaningful. Temperature, queue depth, memory in use.
A histogram records a distribution as a set of cumulative bucket counters. It exists because percentiles cannot be aggregated across sources, which is developed in the aggregation section, and it is the only way to get a correct percentile over a fleet.
A summary computes percentiles at the source. Cheaper to query, and its values cannot be combined across instances at all, which is almost always the wrong trade.