OpenTSDB is a time-series database that stores nothing itself. It is a fleet of stateless Java daemons, called TSDs, that turn metric data points into carefully shaped HBase rows and turn queries back into scans over those rows. Everything durable, including the dictionary that maps metric and tag names to compact IDs, lives in HBase. That split is the whole design: HBase brings horizontal scale, replication and compaction, and OpenTSDB brings a row-key layout that makes time-range reads for one metric a short contiguous scan.
This article explains that layout byte by byte, then shows how data gets in, how queries are answered, where the sharp edges are and how to size a deployment, so you can read its tables, predict its load and fix its most common failures.
The data model
A data point in OpenTSDB has four parts: a metric name, a timestamp, a numeric value and a set of tags. The metric names what is measured, for example sys.cpu.user. Tags are key/value pairs that say where it was measured, such as host=web01 and cpu=0. A time series is one metric plus one exact combination of tag values, so sys.cpu.user{host=web01,cpu=0} and sys.cpu.user{host=web01,cpu=1} are different series. Every data point needs at least one tag, and the number of tags per point is capped by tsd.storage.max_tags, which defaults to 8.
Values are integers or floats; there are no strings. Timestamps are Unix epoch seconds or milliseconds. Queries name a metric, filter on tags and aggregate whatever series match, so your tag design decides both cost and query power.
Architecture
The TSD speaks a line-oriented telnet-style protocol and an HTTP JSON API on the same port, conventionally 4242. With no durable state, TSDs scale by adding instances behind a load balancer and upgrade by rolling restarts; their UID cache and compaction queue are rebuilt or harmlessly lost on restart.
Three HBase tables matter. tsdb holds the data points in a single column family named t. tsdb-uid holds the bidirectional dictionary between strings and UIDs. tsdb-meta is an optional index of time series used for search and metadata. Tree and rollup tables exist for optional features and can be ignored until you need them.
The UID layer
Storing the string sys.cpu.user in every row key would waste space and break fixed-width parsing, so OpenTSDB assigns every metric name, tag key and tag value a unique ID. By default each UID is 3 bytes, which allows 16,777,215 values per kind. The tsdb-uid table stores the mapping in both directions: rows keyed by the string carry the UID in the id family, and rows keyed by the UID carry the string in the name family, each with a qualifier of metrics, tagk or tagv.
Metric UIDs are not created on demand unless you enable tsd.core.auto_create_metrics (default false), so a put for an unknown metric is rejected until you run tsdb mkmetric sys.cpu.user, and a collector typo cannot create a metric. Tag keys and values are assigned as they appear, which is where cardinality problems start: put a request ID or user ID in a tag value and every request burns a UID forever. The tag value space is shared across all metrics, so one bad collector can exhaust it for everyone.
Row keys and qualifiers, byte by byte
The data table row key is [salt]<metric_uid><timestamp><tagk1><tagv1>...<tagkN><tagvN>. The timestamp is 4 bytes and is the base time, normalised down to the hour, not the time of the point. Tag pairs are sorted by tag key UID so the same series always yields the same key. With two tags and no salt the key is 3 + 4 + 2 x (3 + 3) = 19 bytes, and one row holds one series for one hour.
The point's position within the hour goes in the column qualifier. For second-resolution points the qualifier is 2 bytes: the first 12 bits are the offset in seconds from the base time and the last 4 bits are flags for the value type (integer or float) and length. Twelve bits give 4,096 values, enough for the 3,600 seconds in an hour. Millisecond points use a 4-byte qualifier whose first 4 bits are all ones, followed by a 22-bit millisecond offset; 22 bits give 4,194,304 values, enough for 3,600,000 milliseconds. Since qualifiers sort in byte order, a row's cells are already in time order.
| Part | Bytes | Holds |
|---|---|---|
| Salt (optional) | 0 by default | hash bucket of the rest of the key |
| Metric UID | 3 | which metric |
| Base timestamp | 4 | epoch seconds rounded down to the hour |
| Tag key + value UIDs | 6 per tag | series identity, sorted by tag key UID |
| Qualifier (seconds) | 2 | 12-bit offset + 4 flag bits |
| Qualifier (milliseconds) | 4 | 0xF prefix + 22-bit offset + flags |
The layout explains the read path. A query for one metric over a day is a scan from metric_uid + start_hour to metric_uid + end_hour, with a row filter on tag UIDs. Everything for one metric is contiguous, which is fast to read and, as the failure modes show, a hazard to write.
Writing data
The telnet protocol is one point per line, useful for quick tests and simple agents:
put sys.cpu.user 1759546800 42.5 host=web01 cpu=0
put sys.cpu.user 1759546810 40.1 host=web01 cpu=0Production writers should use POST /api/put with batches. A success returns HTTP 204 with no body; if any point fails the TSD returns 400, and ?summary or ?details puts the counts and reasons in the body. A minimal batching writer in Python:
import time, requests
TSD = "http://tsd.internal:4242/api/put?summary"
def flush(points):
if not points:
return
r = requests.post(TSD, json=points, timeout=10)
if r.status_code == 400:
# Unknown metric, too many tags, bad value: the summary body says how many and why.
raise RuntimeError(f"points rejected: {r.text[:500]}")
r.raise_for_status()
batch = []
now = int(time.time())
for cpu, value in enumerate(read_cpu_percentages()): # your collector
batch.append({"metric": "sys.cpu.user", "timestamp": now,
"value": value, "tags": {"host": "web01", "cpu": str(cpu)}})
if len(batch) >= 500:
flush(batch); batch = []
flush(batch)Inside the TSD, a put resolves names to UIDs (from cache, or from tsdb-uid), builds the row key and qualifier, and issues an asynchronous HBase put through the asynchbase client. Writes are buffered and flushed on an interval set by tsd.storage.flush_interval (1,000 ms by default), so a crashed TSD can lose up to that much buffered data. Keep collectors retrying on failure, and treat a 400 as an alert, not a log line.
Compaction, appends and salting
Writing one HBase cell per point is expensive to store, because HBase repeats the full row key, family, qualifier and timestamp for every cell. OpenTSDB therefore runs its own compaction, unrelated to HBase compaction. With tsd.storage.enable_compaction on, the default, the TSD remembers each row it wrote; once the hour is over it reads the row, concatenates all qualifiers and values into a single cell, writes that cell and deletes the originals. Storage drops sharply, but each row is now written, read, rewritten and deleted, so a busy cluster pays a steady read and write load right after every hour boundary.
OpenTSDB 2.2 added appends as an alternative, controlled by tsd.storage.enable_appends (off by default). Each point is appended to a single column with qualifier 0x050000 as offset/value pairs, so no hourly rewrite is needed. The cost moves to the RegionServer: an HBase append is a read-modify-write under a row lock, so appends trade the hourly spike for constant extra work per point. Pick one deliberately and measure RegionServer CPU and latency before and after.
Salting, also from 2.2, prefixes the key with tsd.storage.salt.width bytes (default 0, meaning off) holding a hash bucket from 0 to tsd.storage.salt.buckets minus one (default 20 buckets). Writes for one metric then spread over that many key ranges, and every query fans out into one scanner per bucket and merges the results. Salt settings are part of the key format: change them on a live table and existing data becomes unreadable by the new configuration. Decide before the first write. The general technique is covered in HBase salting.
Queries, downsampling and interpolation
Queries go to /api/query. The compact GET form packs everything into one m parameter: aggregator, optional downsampler, optional rate, metric and tag filters.
GET /api/query?start=1h-ago&m=sum:1m-avg:rate:sys.cpu.user{host=web*}
POST /api/query
{
"start": "24h-ago",
"queries": [{
"metric": "http.requests",
"aggregator": "sum",
"downsample": "5m-sum-zero",
"rate": true,
"rateOptions": {"counter": true, "resetValue": 1000000},
"filters": [{"type": "wildcard", "tagk": "dc", "filter": "*", "groupBy": true}]
}]
}The query guide gives the order of operations: filtering, grouping, downsampling, interpolation, aggregation, rate conversion, then functions and expressions. Counters such as request totals should always be read with a counter rate so that a process restart, which resets the counter, does not show as a huge negative spike.
The aggregation step hides the most important behaviour in OpenTSDB. The sum, avg, dev, min, max and percentile aggregators use linear interpolation: when one series has a point at 10:00:10 and another only at 10:00:00 and 10:00:20, the second series is assigned an estimated value at 10:00:10 so the two can be combined. zimsum substitutes zero for missing points instead, and mimmin and mimmax ignore them. Downsampling first is the usual cure, because it snaps every series onto the same bucket timestamps. Since 2.2 the downsampler takes a fill policy, as in 1m-avg-nan: none (the default, interpolate), nan, null or zero. Version 2.3 added calendar intervals such as 1dc-sum and the 0all-sum form that reduces a whole range to one value.
Worked example: sizing a fleet
Suppose 10,000 hosts each report 50 metrics every 10 seconds with two tags, host and dc. That is 500,000 series and 50,000 points per second. Each series writes 360 points per hour into one row, so the cluster creates 500,000 rows per hour, 12 million per day.
Before compaction, each point is one HBase cell of roughly 45 to 50 bytes: about 20 bytes of row key plus fixed per-cell overhead (lengths, family, 2-byte qualifier, 8-byte HBase timestamp, type) and a value of 1 to 8 bytes. At 50,000 points per second that is about 2.4 MB/s into the WAL and memstores. After TSD compaction, a row becomes one cell holding 720 bytes of qualifiers plus the values, around 2 KB, instead of about 17 KB of separate cells, a reduction of roughly eight times before HFile block compression. Twelve million compacted rows at about 2 KB is around 25 GB per day of logical data, multiplied by HDFS replication. These figures are approximate; measure your own value widths and tag counts.
Now the risk. Without salt, every write for one metric lands at the tail of one region, so with 50 metrics at most 50 regions take all writes, however many RegionServers you own. With salt.width=1 and 20 buckets that grows to as many as 1,000 key ranges; pre-split the table so they start on different servers. A one-metric dashboard query now opens 20 scanners, a trade that write-heavy monitoring almost always wins.
Failure modes
- Write hotspot. One RegionServer at high CPU and queue depth while others idle, usually an unsalted table or a newly created table that was never pre-split. See HBase hotspotting for diagnosis.
- UID exhaustion or bloat. High-cardinality tag values grow
tsdb-uidwithout bound and eventually hit the 3-byte ceiling. Ban IDs, URLs with parameters and timestamps as tag values; review new tag keys in code review. - Compaction storms. Read and write latency climbs just after each hour as TSDs compact the previous hour's rows. Spread load over more TSDs, consider appends, and watch HBase compaction too; HBase compaction covers the storage side.
- Misleading sums. Gaps in one series are filled by interpolation, so a sum over hosts can show traffic from a host that stopped reporting. Downsample with an explicit fill policy or use
zimsum. - Duplicate points. Two writes for the same series and second with different values make the compacted row ambiguous and queries fail.
tsd.storage.fix_duplicates(2.1+, off by default) keeps the latest value; fix the duplicate writer anyway, and usetsdb fsckto repair existing rows. - Expensive wildcard queries. A filter on a tag that matches millions of series reads every matching row and holds results in TSD memory. Set
tsd.query.timeout(0, meaning no limit, by default) and give dashboards their own TSDs so ad-hoc queries cannot starve ingest.
Trade-offs
OpenTSDB's strengths are the ones HBase gives it: horizontal scale to billions of points, cheap long retention and no single storage node to outgrow. If you already run HBase and Hadoop well, adding TSDs is a small step. The costs are operational weight, a query language with fewer functions than modern systems, write-time UID management and interpolation semantics that surprise new users.
Prometheus is simpler for service monitoring at modest scale and has a richer query language, but needs extra components for long retention; see Prometheus at scale. Choose OpenTSDB when you have an HBase platform team and need years of raw data, and check the project's current release activity before starting anything new on it.
What to do next
- Run
scan 'tsdb-uid', {LIMIT => 20}in the HBase shell and decode a few rows so the UID layout is concrete. - Count tag values per tag key and flag any key with unbounded growth.
- Check
tsd.storage.salt.widthand the region count oftsdb; if writes concentrate on a few regions, plan a salted, pre-split table and a migration. - Decide between TSD compaction and appends by measuring RegionServer latency around the hour mark.
- Switch dashboards that sum across hosts to downsampled queries with an explicit fill policy.
- Use counter rate options for every monotonically increasing metric.
- Set a query timeout and separate read TSDs from write TSDs.
- Alert on rejected puts (HTTP 400) and log the
?summarybody with the reason.