HBase is a common landing zone for event streams: device telemetry, clickstreams, ledger entries, feature values for online models. The stream arrives continuously from Kafka or a similar log, and HBase serves point reads and short scans on the result within milliseconds. The write path can absorb this load; HBase write performance covers the server side and HBase client libraries covers BufferedMutator mechanics.
This article covers the part neither does: the semantics of the pipeline between the log and the table. When a worker crashes, what gets written twice? Can a replay corrupt a counter? What happens when events arrive out of order, or a delete arrives before the data it deletes? How should a consumer react when RegionServers push back? The answers mostly come down to row keys and timestamps, and to the order in which you flush and commit.
Delivery semantics and why timestamps matter
A consumer reads records, writes them to HBase, and records how far it got by committing offsets. A crash can happen between any two of those steps, which gives two choices:
- Commit before writing and a crash loses the records in flight. That is at-most-once delivery, rarely acceptable for data you keep.
- Write, flush, then commit and a crash replays everything since the last commit. That is at-least-once delivery, and it is the only sane default.
HBase has no transaction that spans a Kafka offset and a set of rows, so exactly-once delivery is not available end to end. What you can build is effectively-once results: at-least-once delivery into writes that are idempotent, so applying a record twice leaves the table in the same state as applying it once.
HBase makes idempotent puts easy, for a reason worth understanding. A cell is addressed by row, column family, qualifier and timestamp. Writing the same value to the same four coordinates twice produces one logical cell: reads see one version, and compaction discards the duplicate. If you let the server assign the timestamp, a replay writes at a new timestamp, so the replay becomes an extra version. With VERSIONS => 1 that looks harmless for that cell, but any time-range query, version history or change-data feed now sees the event twice. So the first rule of streaming ingest is: derive both the row key and the timestamp from the event, never from the clock of the writer.
Row keys for streams
Stream events usually carry an entity ID and an event time, and the natural row key, entity then time, has two problems. If many writers append the current time for a small set of entities, or if the key starts with a time, every write lands on the last region: the classic hotspot. And if two events for one entity share a timestamp at millisecond resolution, the second overwrites the first.
A durable layout for per-entity time series looks like this:
row key = salt(entity_id) | entity_id | (Long.MAX_VALUE - event_time_ms)
qualifier = event_type [+ ':' + event_seq] # disambiguates same-millisecond events
timestamp = event_time_ms # replays hit the same cell
value = encoded payload
salt(entity_id) = (murmur3(entity_id) mod N_BUCKETS) as one byte, N fixed at table creationThe salt prefix spreads entities across regions. Because it is a function of the entity, all of one entity's rows stay contiguous and a per-entity scan still touches one region; see salting row keys for the read-side cost of salting. The reversed time puts the newest event first, so "latest N events" is a short scan. Pre-split the table on the salt boundaries so ingest starts spread out instead of waiting for splits.
If events can collide on every field you have, add a deterministic sequence from the source, such as the Kafka partition and offset, rather than a random UUID. A random suffix makes every replay a new row, which is exactly the duplicate you were avoiding.
The flush-then-commit loop
Here is the core loop for a Java worker using the plain Kafka consumer and HBase's BufferedMutator. The ordering is the whole point: write, flush, check for failures, then commit.
AtomicReference<Exception> failure = new AtomicReference<>();
BufferedMutatorParams params = new BufferedMutatorParams(TableName.valueOf("telemetry"))
.writeBufferSize(8L * 1024 * 1024)
.listener((RetriesExhaustedWithDetailsException e, BufferedMutator m) -> {
failure.compareAndSet(null, e); // record; decide in the main loop
});
try (Connection conn = ConnectionFactory.createConnection(conf);
BufferedMutator mutator = conn.getBufferedMutator(params)) {
while (running) {
ConsumerRecords<byte[], byte[]> batch = consumer.poll(Duration.ofMillis(200));
for (ConsumerRecord<byte[], byte[]> rec : batch) {
Event ev = Event.decode(rec.value());
Put put = new Put(rowKey(ev));
put.addColumn(CF, qualifier(ev), ev.eventTimeMs(), ev.payload()); // explicit ts
mutator.mutate(put);
}
if (!batch.isEmpty()) {
mutator.flush(); // blocks until buffered mutations are acked
Exception e = failure.getAndSet(null);
if (e != null) {
throw new IngestException("rows failed after retries", e); // no commit: replay
}
consumer.commitSync(); // offsets of everything polled so far
}
}
}Three choices in that loop deserve comment. The listener only records the failure, because the background flush thread is the wrong place to decide whether to commit. On failure the worker exits without committing, and the restarted worker replays from the last good offset; idempotent puts make that safe. The RetriesExhaustedWithDetailsException lists which rows and which servers failed, so log a sample of them before exiting. A malformed record that will never succeed should go to a dead-letter topic before mutate, not crash the worker forever.
Flushing on every poll keeps the replay window small but caps throughput when polls are small. A common compromise flushes and commits when the buffer passes a size or a record count, or after a fixed interval such as one second, whichever comes first. The replay window is then bounded by that interval.
One caution: in this simple loop, flush() and a full-buffer mutate() block the polling thread. Bound the worst-case flush below max.poll.interval.ms with hbase.client.operation.timeout and a modest retry count, or, at higher volume, flush off the poll thread and pause partitions until it succeeds, as the backpressure section describes.
Counters: the replay trap
HBase's Increment and Append operations are the most common way streaming ingest silently corrupts data. An increment is a read-modify-write on the server. A replay applies it again, and the counter is now wrong with no trace of why. HBase does attach nonces to client retries of increments and appends, so the server can reject a retried duplicate within a time window, but a nonce does not survive a worker restart: a replay from Kafka is a brand-new operation.
Three designs avoid that:
- Write facts, aggregate on read or in batch. Store one idempotent cell per event as above, and compute counts with a scan or a scheduled job. Simple and exactly right; reads cost more.
- Pre-aggregate in the stream processor, then write absolute values. A stateful processor such as Flink or Kafka Streams keeps the running count with its own checkpointed state and writes the current total as a Put with a deterministic timestamp, such as the window end. A replay rewrites the same total.
- Deduplicate with a marker.
checkAndMutatewrites the increment only if a marker cell for the event ID is absent, and sets the marker in the same atomic row operation (Increment inside a conditional row mutation needs a recent HBase 2.x release; check yours). This works only when the marker and the counter share a row, costs a read per event, and needs a TTL on markers so they do not grow forever. Use it sparingly.
Option 2 is usually the right answer for dashboards and feature stores; option 1 for audit data where every event matters.
Late events, deletes and TTL
Streams are rarely in order across partitions, and HBase's version rules interact with lateness in ways that surprise people. With event-time timestamps:
- A late put with an older timestamp does not overwrite a newer one. Reads return the highest timestamp, so for a "current value" column, last-writer-wins by event time comes for free. With
VERSIONS => 1, the older cell is dropped at compaction. - A delete masks everything at or before its timestamp, including puts that arrive later. A
Deletefor a row at time T writes a tombstone; a put at T-5 that arrives after it is invisible on read, and is discarded at the next major compaction. If deletes in your stream mean "delete what exists now", that is correct. If a late "create" event can follow a delete legitimately, you need a different model, such as soft-delete columns. - A delete with a future timestamp is a trap. A producer whose clock is ahead writes tombstones that mask legitimate writes until the clock catches up. Validate event times at ingest and reject or clamp anything beyond a small skew window.
- TTL uses cell timestamps. If you write event time and the column family has a TTL, a backfill of year-old events is expired on arrival. Size TTLs with replays and backfills in mind; see TTL and versions.
Backpressure and durability
When a region's MemStores reach their blocking limit, or a RegionServer's call queue fills, the server rejects writes with RegionTooBusyException or CallQueueTooBigException. The client retries with backoff, governed by hbase.client.retries.number and hbase.client.pause, and a BufferedMutator eventually blocks mutate() when its buffer is full. That blocking is your signal.
The dangerous reaction is to stop calling poll(). A Kafka consumer that does not poll within max.poll.interval.ms is evicted from the group, its partitions move to another worker, which hits the same overloaded regions, and a rebalance storm follows. The right reaction is consumer.pause(partitions): keep polling to stay in the group, fetch nothing, and resume() when flushes succeed again. Pair it with server-side limits so one noisy pipeline cannot starve others; see quotas and throttling.
Durability is the other dial. Put.setDurability accepts SYNC_WAL (the usual default), ASYNC_WAL and SKIP_WAL. Asynchronous WAL raises throughput but can lose recently acknowledged writes if a RegionServer dies. In a replayable pipeline that loss is recoverable only if the replay window covers it, and committing offsets after an acknowledged flush means it usually does not. Use SYNC_WAL unless you also delay offset commits, and never SKIP_WAL for data you cannot regenerate. For one-off history loads, skip the write path entirely with bulk load.
Flink and Spark sinks
If the stream is already processed in Flink, its HBase SQL connector is a sink you configure rather than code. According to the Flink documentation, the connector is hbase-2.2, it always works in upsert mode on the primary key declared in the DDL, and it buffers writes:
CREATE TABLE telemetry_hbase (
rowkey STRING,
d ROW<temp DOUBLE, status STRING>,
PRIMARY KEY (rowkey) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
'table-name' = 'telemetry',
'zookeeper.quorum' = 'zk1:2181,zk2:2181,zk3:2181',
'sink.buffer-flush.max-size' = '4mb', -- default 2mb
'sink.buffer-flush.max-rows' = '2000', -- default 1000
'sink.buffer-flush.interval' = '1s' -- default 1s
);The documentation states the upsert behaviour but does not promise exactly-once delivery, so design as for any at-least-once sink: deterministic row keys, and either a deterministic version or a value that is safe to rewrite. The documented options include nothing for the cell timestamp, so expect server-assigned time; replays then add versions, and you should keep VERSIONS => 1 on these tables or deduplicate downstream. Spark Structured Streaming jobs follow the same pattern through foreachBatch, where the batch ID gives you a deterministic token for idempotence.
Worked example: device telemetry
A fleet of 400,000 devices sends a 300-byte reading every two seconds: 200,000 events per second, about 60 MB/s before WAL and replication. The topic has 64 partitions and 16 workers. The table is pre-split into 64 regions on a one-byte salt with 64 buckets, spread over eight RegionServers.
Each worker handles four partitions, about 12,500 events per second, and flushes every 4 MB or second. A crash replays at most about one second of its partitions. Rows follow the salt-entity-reversed-time layout with event time as the timestamp, so replays are invisible. A per-device hourly maximum is computed in Flink and written as absolute values. When a compaction storm pushed two RegionServers into blocking during a test, the workers on those salts paused their partitions; consumer lag rose for four minutes and drained, and no worker left the group.
Failure modes
| Symptom | Likely cause | Fix |
|---|---|---|
| Duplicate versions after a restart | Server-assigned timestamps | Set the timestamp from the event |
| Counters drift upward over weeks | Increments replayed or retried | Write absolute values or facts |
| Rows written but invisible | Tombstone at a later or future timestamp | Validate clocks; soft deletes |
| One RegionServer at 100 percent | Time-leading key or too few salt buckets | Salt, pre-split, rebalance |
| Rebalance storms under load | Worker stopped polling while blocked | pause and resume partitions |
| Data missing after RegionServer crash | ASYNC_WAL or SKIP_WAL with offsets already committed | SYNC_WAL, or delay commits |
| Backfill data vanishes | TTL applied to old event times | Raise TTL or load to a separate table |
What to do next
- Check every writer: are both row key and timestamp derived from the event? Fix any that use writer time or random suffixes.
- Confirm offsets are committed only after a successful
flush(), with listener failures checked first. - Find every
IncrementandAppendon a replayable path and replace it with absolute writes or facts. - Add event-time validation that rejects timestamps beyond a small future skew, especially on deletes.
- Implement pause and resume on backpressure, and alert on consumer lag and rebalance count.
- Run a chaos test: kill a worker mid-batch and a RegionServer under load, then compare row counts and checksums against the source.