Event tracking means recording what users, devices or orders did, one event at a time, and reading it back by entity and by time: the last fifty actions of a user, every login failure for an account in the past day, every checkout event across the site in the last hour. HBase fits this shape well. Writes are appends that land in a log-structured store, rows are sorted by key so a prefix scan returns one entity's events in order, and column-family TTLs expire old data without a deletion job.

The fit depends almost entirely on the row key, and a key that serves one query can make another impossible or create a hot spot. This article designs an event store from its queries: a timeline table keyed for per-entity reads, an index table for per-type reads, the timestamp choice that decides whether TTL works, retention in tiers, a Java writer and reader, a sizing example, and the failure modes that show up in production. Ingest-path mechanics such as replay and backpressure are covered in a separate article and linked rather than repeated.

Start from the queries

Start from the reads, because HBase gives you one sorted index per table, the row key, and nothing else for free. A typical event product needs these:

QueryFrequencyServed by
Last N events for one userVery high (profile pages, support tools)Timeline table, prefix scan, newest first
One user's events in a time windowHighTimeline table, bounded range scan
All events of one type in a time windowMedium (dashboards, alerts)Index table, one scan per bucket
Count of events per user per dayHigh, but tolerates lagRollup table of counters, or batch aggregation
Full history exportRareSnapshot plus a batch job, never live scans

Anything not in this table belongs in an analytical store fed from the same stream. A tall layout, one row per event, suits event data better than a wide one with one row per user and one column per event: tall rows keep each row small, split cleanly, and avoid the ceilings described in HBase wide row limits. The general method behind this table is in HBase schema design, in depth.

Row keys for the timeline

The timeline key is four fixed-width fields. A 2-byte salt derived from a hash of the user ID spreads users across regions. The 16-byte user ID groups one user's events contiguously. An 8-byte reverse timestamp, the maximum long minus the event time in milliseconds, sorts newest first, so the last N events are the first N rows of the prefix. A 16-byte event ID ends the key, so two events at the same millisecond do not collide and a retried write lands on the same row instead of creating a duplicate.

Fixed widths matter: with variable-length IDs, one user's key can be a prefix of another's, and a prefix scan returns both. The salt is a hash of the user, not a random value, so every read can recompute it and a user's events stay in one place; why this differs from random salting is explained in HBase salting row keys. A monotonically increasing key, such as a timestamp first, would send every write to the last region; the hash prefix is what spreads load.

Event tracking with two tables: a per-user timeline and a per-type indexEvent producersapps, servicesLog / queueordered, replayableIngest writerBufferedMutatorevents (timeline)TTL 90 daysevents_by_type (index)TTL 30 days1st2ndTimeline row key (42 bytes, fixed width)salt2 bytesuser_id16 bytes (UUID)reverse event time8 bytes: MAX - millisevent_id16 bytes (UUID)hash(user)groups one usernewest firstmakes retries idempotentIndex row keybucket1 byteevent_type4-byte type codereverse event time8 bytesuser_id + event_idpoints back to the timeline rowReads: last N for a user = one prefix scan; type in a window = one scan per bucket, merged
A timeline table keyed by salt, user, reverse time and event ID, plus an index table keyed by type for cross-user reads. The writer fills the timeline first and the index second.

Two clocks: event time and cell timestamps

An HBase cell has a timestamp separate from anything in the key, and that timestamp is what TTL compares against. A column family with a TTL of 90 days drops a cell when its timestamp is more than 90 days older than the current time. If you write cells with the event time as their timestamp, a late event that arrives 91 days after it happened expires on arrival; it is written, acknowledged and invisible. More subtly, an event replayed from a backlog a week late gets a week less retention than its neighbours.

The design here separates the two clocks. Event time goes into the row key, where it controls ordering and range scans. The cell timestamp is left to the server, which assigns the write time, so TTL means time since ingest and retention is predictable. Store the event time again as a column value if consumers need it with full precision.

Deletes have a related trap. A delete marker hides every cell in its scope with a timestamp at or before the marker's. If you delete a user's events, for example for an erasure request, and the pipeline later re-puts an older event with an explicit old timestamp, the new put is masked and reads miss it. It stays masked until a major compaction removes the marker. Server-assigned timestamps avoid the problem for ordinary writes; for erasure, delete every row in the user's key range and record the request so replays skip that user. Per-family TTL and versions are covered in HBase TTL and MAX_VERSIONS.

Writing events

The writer builds keys, sends puts through a BufferedMutator for batching, and writes the timeline row before the index row. Events come from a durable, replayable log, so a crash between the two writes is repaired by replay, and the deterministic key makes replay overwrite rather than duplicate.

static final byte[] D = Bytes.toBytes("d");
static final byte[] I = Bytes.toBytes("i");

static byte[] salt(byte[] userId) {
    return Arrays.copyOf(MD5.digest(userId), 2);           // stable per user
}

static byte[] timelineKey(byte[] userId, long eventMillis, byte[] eventId) {
    return ByteBuffer.allocate(42)
        .put(salt(userId)).put(userId)                         // 2 + 16
        .putLong(Long.MAX_VALUE - eventMillis)                 // 8, newest first
        .put(eventId).array();                                 // 16
}

static byte[] indexKey(int typeCode, long eventMillis, byte[] userId, byte[] eventId) {
    byte bucket = (byte) ((userId[15] & 0xff) % INDEX_BUCKETS);
    return ByteBuffer.allocate(1 + 4 + 8 + 16 + 16)
        .put(bucket).putInt(typeCode)
        .putLong(Long.MAX_VALUE - eventMillis)
        .put(userId).put(eventId).array();
}

void write(Event e, BufferedMutator timeline, BufferedMutator index) throws IOException {
    Put t = new Put(timelineKey(e.userId, e.millis, e.eventId));
    t.addColumn(D, Bytes.toBytes("type"), Bytes.toBytes(e.typeCode));
    t.addColumn(D, Bytes.toBytes("ts"), Bytes.toBytes(e.millis));
    t.addColumn(D, Bytes.toBytes("body"), e.payload);           // compact, e.g. protobuf
    timeline.mutate(t);

    Put x = new Put(indexKey(e.typeCode, e.millis, e.userId, e.eventId));
    x.addColumn(I, Bytes.toBytes("s"), e.summary);              // small: enough for a list view
    index.mutate(x);
}
// After each batch from the log: timeline.flush(); index.flush(); then commit the log offset.

No explicit timestamp is passed, so the server assigns write time and TTL behaves as described. Flushing both mutators before committing the log offset is what makes the pipeline at-least-once rather than at-most-once. The event ID in the key turns at-least-once into effectively once for the event tables. Counters are different, because an increment replayed twice counts twice; the patterns for that, and for backpressure and late data, are in HBase streaming ingest patterns.

Reading timelines and windows

Reads are range scans with explicit bounds. The last N events for a user scan from the user's prefix and stop after N rows. A time window converts its bounds through the reverse timestamp: the newer bound becomes the start row and the older bound the stop row.

List<Result> lastN(Table t, byte[] userId, int n) throws IOException {
    byte[] prefix = ByteBuffer.allocate(18).put(salt(userId)).put(userId).array();
    Scan s = new Scan()
        .withStartRow(prefix)
        .withStopRow(Bytes.unsignedCopyAndIncrement(prefix), false)
        .setLimit(n).setCaching(n)
        .addFamily(D);
    try (ResultScanner rs = t.getScanner(s)) {
        List<Result> out = new ArrayList<>();
        for (Result r : rs) out.add(r);
        return out;
    }
}

Scan window(byte[] userId, long fromMillis, long toMillis) {
    byte[] p = ByteBuffer.allocate(18).put(salt(userId)).put(userId).array();
    byte[] start = ByteBuffer.allocate(26).put(p).putLong(Long.MAX_VALUE - toMillis).array();
    // +1 so events at exactly fromMillis (any event ID) fall before the exclusive stop row
    byte[] stop  = ByteBuffer.allocate(26).put(p).putLong(Long.MAX_VALUE - fromMillis + 1).array();
    return new Scan().withStartRow(start).withStopRow(stop, false).addFamily(D);
}

For the index table, one query becomes one scan per bucket with the same type and time bounds, merged on the client by reverse timestamp. With 16 buckets that is 16 parallel scans, each small. Buckets keep a hot event type, such as page views, from landing on a single region, at the cost of fan-out on read. Index rows carry a short summary so list views need no second lookup; detail views fetch the timeline row by its key. For how each scan RPC is sized and why limit and caching matter, see HBase scans, in depth.

Retention in tiers

Retention is set per column family and per table, which allows tiers. Raw timeline events keep 90 days for support and user-facing history. The index keeps 30 days, because dashboards rarely look further back and the index is the larger write load per event. A rollup table keyed by salt, user and day holds counters per event type, kept for two years at a tiny fraction of the raw size. Anything older lives in an archive built from snapshots or the original log, not in the live tables.

TTL expiry is lazy: expired cells vanish from reads at once but occupy disk until a compaction rewrites their files, so plan capacity for TTL plus the compaction lag. Because old data is spread across every region rather than concentrated in a few, every region keeps compacting expired cells, so compaction throughput belongs in the capacity plan.

Worked example: 200 million events a day

Size a store for 200 million events a day with an average stored size of 400 bytes per timeline row, including key and column overhead, and an index row of 120 bytes.

QuantityArithmeticResult
Average write rate200,000,000 / 86,400 sabout 2,300 events/s, 4,600 puts/s with index
Peak write rate5x average for evening peaksabout 23,000 puts/s
Timeline per day200 M x 400 B80 GB/day before compression
Timeline retained80 GB x 90 days7.2 TB logical, 21.6 TB with 3 HDFS replicas
Index retained200 M x 120 B x 30 days720 GB logical, 2.2 TB replicated

Compression typically shrinks event payloads severalfold, but measure it on your own data before relying on a ratio. With a target region size of around 10 to 20 GB after compression, the timeline needs somewhere between roughly one and five hundred regions at steady state, depending on the compression ratio you measure. Pre-split the table on salt boundaries, for example 256 splits at evenly spaced 2-byte prefixes, so ingest spreads across servers from day one instead of waiting for splits. A peak of 23,000 small puts a second is modest for a cluster of a handful of RegionServers when writes are batched, but check WAL throughput and memstore flush rates under replayed backlogs, which can run far above live peaks.

Failure modes

  • Hot region from a celebrity entity. One user, device or bot emits millions of events and its contiguous prefix lands on one region. Detect top prefixes, rate-limit at ingest, or add a sub-bucket for known heavy entities.
  • Late events expire on arrival. Cell timestamps were set to event time. Let the server assign timestamps and keep event time in the key.
  • Re-put events invisible after an erasure delete. Delete markers mask older-timestamped puts until major compaction. Skip erased users at ingest.
  • Index drift. The writer crashed between timeline and index writes and the log offset was committed anyway. Commit offsets only after both flushes, and run a periodic reconcile that compares index rows with timeline rows for a sample of users.
  • Unbounded scans. A tool scans a user's whole history without a limit, and a heavy user times out the RPC. Always set a limit or a time bound.
  • Disk grows past the TTL plan. Major compactions are disabled or throttled, so expired cells stay on disk. Monitor store file size against expected retained data.

Trade-offs

ChoiceGainsCosts
Tall rows (one per event)Small rows, clean splits, simple TTLMore keys; per-row overhead
Hash salt on the userSpread writes, single-scan user readsCross-user time scans need the index
Bucketed index tableType queries without full scansDouble writes; fan-out reads; eventual consistency
Server-assigned timestampsPredictable TTL; no masking trapsCannot use timestamps for event-time versioning
Rollup countersCheap aggregate readsNot idempotent under replay without extra design

If most queries are aggregations across many users, HBase is the wrong primary store and a columnar analytical engine fed from the same log is the better fit; keep HBase for the per-entity lookups it serves in milliseconds. When a single entity or event type does turn hot, diagnose it from region and key metrics before changing the key.

What to do next

  1. Write down your event queries with frequency and latency targets before choosing any key.
  2. Use a fixed-width timeline key: hash salt, entity ID, reverse event time, event ID.
  3. Leave cell timestamps to the server and store event time in the key and as a column.
  4. Write the timeline row first and the index second, and commit log offsets only after both flushes.
  5. Set TTL per table in tiers, and add compaction lag to the disk capacity plan.
  6. Pre-split on salt boundaries and load-test with a replayed backlog, not only live rates.
  7. Put limits or time bounds on every scan, and alert on top prefixes by write rate.
  8. Schedule an index reconcile job and a check that erased users are skipped at ingest.
Key takeaway: An HBase event store works when its row keys follow its queries: a salted, fixed-width timeline key of entity, reverse event time and event ID for per-entity reads, and a bucketed index table for per-type reads. Keep event time in the key and let the server assign cell timestamps so TTL measures time since ingest, write timeline before index and commit log offsets after both, tier retention by table, and bound every scan.