An IoT platform turns a device fleet into an endless stream of small, timestamped readings. The writes never stop, almost nobody updates old data, and the reads come in a few predictable shapes: what is this device doing right now, what did it do over the last day, and how has it trended over a year. HBase fits that profile well because it is a log-structured store that turns random writes into sequential ones and serves sorted range scans cheaply, but only if the table design matches the read shapes.
This article designs the storage layer for a device fleet from first principles. It covers the four tables most fleets need, hotspot-free row keys, event-time cell timestamps, clock skew and late uploads, retention and compaction, sizing for 200,000 devices, and production failure modes.
Start from the read patterns
Start with the access patterns, because in HBase the row key is the only index. A typical fleet has four:
- Current state. A dashboard or mobile app asks for the last known value of every metric for one device. This must be a single Get, measured in milliseconds, at high request rates.
- Recent history. A support engineer opens one device and wants its raw readings for the last few hours or days. This is a short range scan over one device's data.
- Long trends. Charts of hourly or daily averages over months. Reading raw data for this is wasteful; it should come from pre-aggregated rows.
- Fleet-wide questions. Which devices on firmware 4.2 overheated last week? This touches every device. HBase can answer it only with a full scan, so plan to serve it elsewhere.
One table cannot make all four cheap, because each wants a different sort order. So write each reading once to Kafka and fan it out into purpose-built tables, trading storage for reads that are all Gets or narrow scans.
Reference architecture
Devices publish to an MQTT broker or an HTTPS endpoint, and a bridge writes every message to Kafka keyed by device ID, so one device's readings stay in order within a partition. An ingest writer consumes Kafka, validates and normalises each reading, and writes two tables: the raw telemetry table and the latest-state table. A separate stream job computes windowed aggregates and writes the rollup table. The device registry (owner, model, firmware, location) changes rarely and is written by the provisioning service.
Kafka is not decoration: devices reconnect in bursts after outages, and HBase regions split, move and flush. A durable buffer absorbs both without dropping readings.
Four tables
| Table | Row key | Columns | Retention |
|---|---|---|---|
| telemetry | salt | device_id | hour_start | one cell per reading, qualifier = seconds into the hour | TTL 30 days |
| latest | device_id | one cell per metric, timestamp = event time | VERSIONS = 1, no TTL |
| rollup_1h | salt | device_id | hour_start | min, max, sum, count per metric | TTL 2 years |
| device | device_id | model, firmware, owner, provisioned_at | forever |
Each table gets its own column family settings. Telemetry and rollups use data block encoding and compression because their keys repeat heavily. The latest table is small and read hot, so it benefits from a generous block cache share and, if the read rate demands it, the IN_MEMORY flag on its family.
Row keys for raw telemetry
The raw telemetry key must satisfy two goals that pull against each other: spread writes across all region servers, and keep one device's readings contiguous so a history query is one scan. A key that starts with a timestamp fails the first goal badly, because every write in the fleet lands at the end of the key space, in one region on one server. A key that starts with the device ID meets both goals as long as device IDs are well distributed. Sequential IDs such as serial numbers are not, so prefix them with a short hash of the ID:
// Java: one key builder, shared by every writer and reader.
static final int BUCKETS = 32;
static byte[] telemetryKey(String deviceId, long eventMillis) {
byte[] dev = Bytes.toBytes(deviceId);
byte salt = (byte) Math.floorMod(stableHash32(dev), BUCKETS); // any fixed 32-bit hash, e.g. Murmur3
long hourStart = eventMillis - Math.floorMod(eventMillis, 3_600_000L);
return Bytes.add(new byte[] { salt }, dev, Bytes.toBytes(hourStart));
}
static byte[] offsetQualifier(long eventMillis) {
int secondsIntoHour = (int) (Math.floorMod(eventMillis, 3_600_000L) / 1000);
return Bytes.toBytes((short) secondsIntoHour); // 2 bytes, sorts in time order
}The salt is a deterministic function of the device ID, never a random number, so a reader can recompute it and issue a single Get or scan. Because the device ID follows the salt, a history query still touches one contiguous range. The hour bucket makes the row wide: a device reporting every 10 seconds produces 360 cells per row, which keeps a day of history to 24 rows instead of 8,640. Use fixed-width device IDs, or a terminator byte, so a scan for device abc cannot match abcd. Pre-split the table on the salt boundaries so all 32 buckets start on different regions instead of waiting for splits.
A common misconception is that wide rows save disk because the row key is stored once. HBase stores the full key (row, family, qualifier, timestamp) with every cell. What makes the repetition cheap is data block encoding such as FAST_DIFF, which stores only the bytes that differ from the previous key in the block. Enable it, and pair it with block compression, on both time-series tables.
The latest-state table
The latest-state table holds one row per device and one cell per metric. The trick that makes it robust is to write each cell with an explicit timestamp equal to the reading's event time and to keep only one version:
create 'latest', {NAME => 's', VERSIONS => 1, BLOOMFILTER => 'ROW', BLOCKCACHE => true}
// Java ingest writer
Put put = new Put(Bytes.toBytes(deviceId));
put.addColumn(S, Bytes.toBytes("temp_c"), reading.eventMillis(), Bytes.toBytes(reading.tempC()));
put.addColumn(S, Bytes.toBytes("batt_v"), reading.eventMillis(), Bytes.toBytes(reading.battV()));
latestMutator.mutate(put);HBase returns the cell with the highest timestamp, not the one written last. If a device buffers readings while offline and uploads them an hour later, those old readings arrive after newer ones, carry older timestamps, and never overwrite the current value. No read-modify-write and no conditional update is needed, so the writer stays idempotent: replaying a Kafka partition after a crash produces exactly the same cells.
The same property is also the table's main risk. A device with a broken real-time clock that reports the year 2031 writes a cell that outranks every honest reading until 2031 arrives. The ingest writer must therefore validate timestamps before using them, which is the next section.
Event time, clock skew and the TTL trap
Every reading has two times: when the device measured it and when the platform received it. Use event time for the cell timestamp and the key, because that is what users query by, but record arrival time too, and apply three rules at ingest:
MAX_FUTURE_MS = 5 * 60 * 1000 # tolerate small clock drift
MAX_AGE_MS = 29 * 24 * 3600 * 1000 # just inside the 30-day telemetry TTL
def classify(event_ms, arrival_ms):
if event_ms > arrival_ms + MAX_FUTURE_MS:
return "quarantine" # bad clock: never let it into 'latest'
if event_ms < arrival_ms - MAX_AGE_MS:
return "archive" # older than TTL: HBase would hide it on arrival
return "accept"The second rule exists because TTL is measured against the cell timestamp, not the write time. A backfill of readings that are 40 days old into a table with a 30-day TTL succeeds without error, yet the cells are invisible to every read immediately and are dropped at the next compaction. Route such data to an archive (object storage, or a table with a longer TTL) instead.
Deletes interact with timestamps in the same way. A Delete without an explicit timestamp writes a tombstone at the current time, masking every cell at or below it, including readings not yet uploaded that carry older event times; when they arrive they stay hidden, and the next major compaction discards them with the tombstone. Conversely, a cell with a future timestamp survives a delete issued now. When a device is decommissioned, prefer to stop writing and let TTL expire its telemetry rather than issuing deletes against live rows.
Rollups, retention and compaction
Trend charts over months should never scan raw readings. A stream job (Spark Structured Streaming or Flink) groups readings by device and hour, computes min, max, sum and count per metric, and writes them to the rollup table. Store sum and count rather than the average, so coarser rollups can be derived exactly by adding.
Write rollups as Puts that overwrite the hour's cells, not as HBase Increments. An Increment is not idempotent, so a job that restarts and replays a window double-counts. A Put of the recomputed total for the window is safe to repeat. Late readings that arrive after a window was emitted are handled the same way: the job re-emits the corrected hour and the Put replaces it.
Time-ordered data with TTL is also a good fit for the date-tiered compaction policy, which groups store files by the age of the data they contain so that old files are compacted rarely and whole files can be dropped when their data expires. It is enabled per table with the hbase.hstore.engine.class property set to org.apache.hadoop.hbase.regionserver.DateTieredStoreEngine. It has a real cost: heavy out-of-order backfill writes old timestamps into new files, which forces old tiers to be recompacted and erodes the benefit. Use it for tables whose writes are mostly in time order, and test it with a realistic backfill before enabling it on production telemetry.
Worked example: sizing 200,000 devices
Take a fleet of 200,000 devices, each sending 8 metrics every 10 seconds. That is 20,000 readings per second. The storage cost depends heavily on whether each metric is its own cell. A cell in an HFile carries roughly 4 bytes of key length, 4 of value length, 2 of row length, the row key, 1 byte of family length, the family, the qualifier, 8 bytes of timestamp and 1 byte of type, then the value.
| Layout | Cells/s | Bytes per cell | Raw MB/s | Raw GB/day |
|---|---|---|---|---|
| One cell per metric (8-byte double) | 160,000 | about 56 | about 9.0 | about 774 |
| One packed cell per reading (64-byte record) | 20,000 | about 112 | about 2.2 | about 194 |
The figures assume a 25-byte row key, a 1-byte family name and a 2-byte qualifier, and they are before block encoding, compression and HDFS replication. Packing all metrics of one reading into a single cell (a fixed binary record or a small protobuf) cuts the cell count eightfold and storage about fourfold, at the cost of reading the whole record to get one metric. For telemetry that is almost always read as a full reading, packing is the better default; keep the latest table one cell per metric because dashboards often want a single value.
With 30 days of retention, the packed layout holds roughly 5.8 TB of raw data before compression, which then multiplies by three on HDFS. Divide by your target region size to get a region count, and spread those regions across enough servers that each handles a few thousand writes per second with headroom for reconnect storms. Load-test the per-server figure on your own hardware.
Failure modes
- Reconnect storm. A regional network outage ends and tens of thousands of devices upload buffered readings at once. Without Kafka in front, client retries pile into region servers and memstores hit their blocking limits. Keep the buffer, cap writer concurrency, and let lag drain.
- Hot device. A misconfigured device reporting every 10 milliseconds turns its row into a hotspot no salt can spread. Enforce per-device rate limits at the broker and alert on per-key write rate.
- Future timestamps. A bad clock pins a stale value in the latest table for years. Quarantine readings beyond the future tolerance, and give support staff a tool that removes a specific cell by timestamp.
- Silent backfill loss. Replays older than the TTL vanish. Classify by age at ingest and alert when the archive path receives data.
- Region count creep. Long retention plus small regions leaves thousands of regions per server and slows recovery. Size regions for the retained volume and revisit when retention changes.
Trade-offs
HBase is a strong choice when the fleet is large, writes are continuous, the team already runs Hadoop-family infrastructure, and reads are per-device. It is a weaker choice when most questions are fleet-wide aggregations, when the team has no HBase operators, or when the data is small enough for a purpose-built time-series database on a single server. Fleet-wide analytics belongs in a columnar store fed from Kafka or from periodic HBase snapshots exported to Parquet, where a scan of one metric across all devices reads only that column.
The four-table design trades storage and write amplification for read predictability, usually the right trade for IoT.
Related reading
For the details this page summarises, see HBase key salting for bucket counts and pre-splitting, streaming ingest into HBase for the flush-then-commit writer loop, TTL and versions for expiry semantics, wide versus tall tables for row-shape trade-offs, and hotspotting for diagnosing uneven regions.
What to do next
- Write down your fleet's four read patterns and the latency each needs before creating any table.
- Build one shared key-builder library with a deterministic salt and use it in every writer and reader.
- Create the latest table with VERSIONS = 1 and write cells with event-time timestamps.
- Add ingest checks that quarantine future timestamps and divert readings older than the TTL.
- Compute rollups in a stream job and write them as idempotent Puts, never Increments.
- Run the sizing arithmetic for your fleet, pre-split on salt boundaries and load-test a reconnect storm.
- Route fleet-wide analytics to a columnar export instead of HBase scans.