An HBase schema is not a set of tables with relationships. It is a choice of byte strings. HBase keeps every table sorted by row key, splits it into regions by key range, and offers exactly one way to find data quickly: by row key or a contiguous range of row keys. Everything the application needs to read efficiently must therefore be encoded, in the right order, into keys.

This page teaches a method: list the reads, turn each read into a key range, then choose the byte layout that makes those ranges contiguous and the writes evenly spread. It works through a complete messaging schema, covers the encoding mistakes that break sort order, and draws the line between what HBase makes atomic and what you have to make safe yourself. Neighbouring topics have their own pages and are linked where they come up: hotspotting and salting, column family configuration, and TTL and versions.

Advertisement

The model you are designing against

Logically, an HBase table is a sorted map: row key to column family to column qualifier to timestamp to value. All five parts are byte arrays except the timestamp, a long. Rows are stored in unsigned lexicographic order of their key bytes. A region holds a contiguous range of rows and is served by one RegionServer; when a region grows past its size threshold (hbase.hregion.max.filesize, 10 GB by default in HBase 2) it splits at a row boundary.

Four consequences drive every design decision. There is one index. A Get by row key or a Scan over a key range is fast; a lookup by anything else is a full scan unless you build another table. Adjacent keys live together. A range read touches one or a few regions, and a burst of writes to adjacent keys lands on one server. A row never splits. Whatever you put in one row lives in one region on one server. Atomicity stops at the row. A mutation to one row is atomic across all its families and columns; nothing spanning rows is, except in narrow cases covered below.

Design from reads: every query becomes a row-key range on one tableQ1 inbox, newest firstrange on user_id prefixQ2 messages in a threadrange on thread prefixQ3 unread countsingle-row Getinboxuser_id | rev_ts | thread_idmessagesbucket | thread_id | msg_sequser_stateuser_id; counters as cellsSorted by unsigned bytes; regions are contiguous key rangesregion [ , 1f)region [1f, 4a)region [4a, b0)region [b0, )Atomic: one rowany families and columns, checkAndMutate, IncrementNot atomic: across rows or tablesindex tables need write order plus repairThere is one index, the row key. Everything else is a scan or a table you maintain.
The worked example: three reads, three tables, each read a contiguous key range or a single row. The atomicity boundary at the bottom decides how the tables are kept consistent.

Step one: write the reads down

Relational design starts from entities. HBase design starts from queries, because a query that no key serves is a full-table scan. Write every read the application performs, with its frequency, latency target and result size, before drawing a single table. The worked example is a messaging service with 50 million users.

ReadFrequencyShape
Q1: a user's inbox, newest threads first, 20 at a timevery highrange of one user's rows, descending time
Q2: messages in a thread, in order, pagedhighrange of one thread's rows, ascending sequence
Q3: a user's unread countvery highone value
Q4: find a message by its global idlow (support tools, deep links)point lookup by a non-key attribute

Writes are one new message: append to the thread, move the thread to the top of each participant's inbox, and bump unread counts. Every read above must be a Get or a bounded Scan on some table. Q4 does not match any natural key, which tells you already that it needs an index table.

Advertisement

Step two: turn each read into a key

Inbox table. Row key user_id | reverse_ts | thread_id. The user id prefix makes Q1 a range; the reverse timestamp, Long.MAX_VALUE - last_message_millis, makes the newest thread sort first, because HBase only scans forward efficiently (reverse scans exist but cost more); the thread id disambiguates threads updated in the same millisecond. Moving a thread to the top is a delete of its old inbox row and a put of a new one, so the user_state row keeps each thread's current timestamp in order to find the old key.

Messages table. Row key bucket | thread_id | msg_seq, one row per message. Thread ids are random, so they spread writes naturally; the one-byte bucket derived from a hash of the thread id makes pre-splitting straightforward. If thread ids were sequential, this is where the salting and hashing techniques in the hotspotting guide would be required. The sequence number is a fixed-width counter within the thread, so Q2 is a forward range scan.

User state table. Row key user_id with one cell per counter and small per-thread pointers. Q3 is a single Get, and the unread count is updated with Increment, which is atomic within the row.

Message index table. Row key message_id with one cell holding the messages-table key. Q4 becomes two Gets. This is an application-maintained index, discussed below.

All four tables use a single column family. Multiple families are justified only in specific cases, and the column family design guide explains when; retention of old messages is a TTL question, covered by the TTL and versions guide.

Byte encoding: where designs quietly break

HBase compares keys as unsigned bytes. Most encoding bugs come from forgetting that. Signed integers. Bytes.toBytes(long) writes big-endian two's complement, so negative numbers begin with a byte of 0x80 or more and sort after all positive numbers. If a key component can be negative, flip the sign bit before writing. Variable-length strings. If a component is a string followed by another component, ab|c and a|bc can collide or sort wrongly unless you either pad to a fixed width or use a separator byte that cannot appear in the string, typically 0x00. Decimal text for numbers. The string "10" sorts before "9"; store numbers in fixed-width binary or zero-padded text. Key length. The row key is stored in every cell, so a 100-byte key on a table with ten columns per row stores that key ten times on disk; keep keys short and prefer fixed-width binary ids to readable strings. The hard maximum is 32,767 bytes, but anything near it is a design error.

import org.apache.hadoop.hbase.util.Bytes;

final class Keys {
    // Inbox: user_id (8) | reverse_ts (8) | thread_id (8) = 24 bytes, fixed width.
    static byte[] inbox(long userId, long lastMsgMillis, long threadId) {
        return Bytes.add(Bytes.toBytes(userId),
                         Bytes.toBytes(Long.MAX_VALUE - lastMsgMillis),
                         Bytes.toBytes(threadId));
    }

    // Messages: bucket (1) | thread_id (8) | seq (8). Bucket spreads threads over pre-split regions.
    static byte[] message(long threadId, long seq) {
        byte bucket = (byte) (Long.hashCode(threadId) & 0x0F);   // 16 buckets
        return Bytes.add(new byte[] {bucket}, Bytes.toBytes(threadId), Bytes.toBytes(seq));
    }

    // Sortable encoding for a component that may be negative: flip the sign bit.
    static byte[] sortableLong(long v) {
        return Bytes.toBytes(v ^ Long.MIN_VALUE);
    }
}

Here the ids and timestamps are non-negative, so plain Bytes.toBytes sorts correctly. The sortable helper is for fields such as scores or offsets that can go below zero.

Tall versus wide

The inbox could instead be one row per user with one column per thread: a wide design. Both work for small data, and the choice is about limits. A wide row cannot split, so a user with a million threads makes one enormous row on one server that strains scans, RPCs and memory, and any single cell larger than the client's hbase.client.keyvalue.maxsize (10 MB by default) is rejected. Paging through columns within a row is possible with column pagination filters but clumsier than a row range. A tall design, one row per item, spreads a heavy user over as many regions as needed and pages with ordinary scans.

Wide rows are the right choice when the set of columns is bounded and always read together, and when you need them to change atomically: a user profile, a small set of counters, a document with a few dozen fields. That is why user_state is wide and the inbox is tall. The rule of thumb: tall for anything that grows without a known bound, wide for bounded sets that must move together.

What is atomic, and what you must make safe

A Put, Delete, Increment or Append on one row is atomic across all its families and columns, and readers see either all of it or none of it. checkAndMutate adds a conditional write on one row, which is how you implement create-if-absent and optimistic concurrency. MultiRowMutationEndpoint can apply mutations to several rows atomically, but only when they are in the same region, which the key design must guarantee, and it is a coprocessor that has to be enabled on the table.

Nothing spans tables. Posting a message writes the messages table, deletes and puts inbox rows for every participant, increments unread counters and writes the index row. Any of those can fail after the others succeed. Design for it: write in an order that leaves recoverable states, make each write idempotent, and repair what is left.

// 1. Source of truth first, create-if-absent so a retry cannot duplicate.
CheckAndMutate create = CheckAndMutate.newBuilder(Keys.message(threadId, seq))
        .ifNotExists(CF, BODY)
        .build(new Put(Keys.message(threadId, seq)).addColumn(CF, BODY, body));
boolean firstAttempt = messages.checkAndMutate(create).isSuccess();   // false on a retry

// 2. Index next: it points at a row that now exists. A plain Put is idempotent.
index.put(new Put(Bytes.toBytes(messageId)).addColumn(CF, REF, Keys.message(threadId, seq)));

// 3. Derived views last. Inbox moves are idempotent and always replayed;
//    the increment is not, so it runs only on the attempt that created the message.
for (long user : participants) {
    moveThreadToTop(user, threadId, sentAtMillis);   // delete old inbox row, put new one
    if (firstAttempt) userState.incrementColumnValue(Bytes.toBytes(user), CF, UNREAD, 1L);
}

With this order, a crash can leave a message with no index row, no inbox entry or a missed counter, but never an index entry pointing at nothing, and a retry of the same post completes the idempotent steps. For posts that are never retried, a background job scans recent messages and re-creates missing index rows, inbox rows and counters, and readers of the index verify the referenced row exists. Counters deserve care: Increment is not idempotent, so the code above skips it on retries, which means a crash between the message write and the increment undercounts; recompute counts periodically from the inbox to correct that. CheckAndMutate as a builder object requires an HBase 2.4 or later client. If you want indexes maintained for you with SQL on top, Apache Phoenix provides global and local secondary indexes, at the cost of another layer to operate.

Reading it back: prefix scans in the HBase 2 client

Q1 and Q2 are prefix scans. In the HBase 2 client, set a start row and a stop row rather than a prefix filter: the region server then seeks directly to the start and stops at the boundary, while a filter alone may read far more rows than it returns. The stop row for a prefix is the prefix with its last byte incremented, carrying over any 0xFF bytes.

static byte[] stopRowFor(byte[] prefix) {
    byte[] stop = Arrays.copyOf(prefix, prefix.length);
    for (int i = stop.length - 1; i >= 0; i--) {
        if (stop[i] != (byte) 0xFF) { stop[i]++; return Arrays.copyOf(stop, i + 1); }
    }
    return HConstants.EMPTY_END_ROW;    // prefix was all 0xFF: scan to the end
}

byte[] prefix = Bytes.toBytes(userId);
Scan inboxPage = new Scan()
        .withStartRow(prefix)
        .withStopRow(stopRowFor(prefix), false)
        .setLimit(20)                    // one page of threads
        .setCaching(20);
try (ResultScanner rs = inbox.getScanner(inboxPage)) {
    for (Result r : rs) render(r);
}

For the next page, start from the last row returned with the start row marked exclusive. The scans guide covers caching, batching and timeouts in more detail.

Failure modes

  • Monotonic leading key. A timestamp or sequence at the front of the key sends all writes to the last region. Lead with a high-cardinality id or a bucket.
  • Unbounded wide rows. One heavy user or device makes a row that cannot split and eventually exceeds the cell or RPC size limits.
  • Sort-order bugs. Negative numbers, decimal text and unseparated strings sort wrongly; they usually surface as missing results at page boundaries, long after launch.
  • Queries nobody designed for. A new feature filters by an attribute not in any key and ships as a full scan with a filter. Add a table or an index first.
  • Assuming cross-row atomicity. Index and view tables drift after partial failures unless writes are ordered, idempotent and repaired.
  • Too many families. Splitting data into families by entity rather than by access pattern multiplies flushes and store files.

What to do next

  1. List every read with its frequency, latency target and result size before creating a table.
  2. For each read, write the key range it will scan; if you cannot, add a table or an index table.
  3. Define keys as fixed-width binary components, check signed values and separators, and unit-test the sort order with edge values.
  4. Choose tall rows for anything unbounded and wide rows only for bounded sets that must change atomically.
  5. Order multi-table writes so failures leave recoverable states, make them idempotent, and schedule a repair job.
  6. Pre-split tables using the known key distribution, and load-test with production-like keys to find hotspots before launch.
Key takeaway: HBase gives you one sorted index, the row key, and one atomic unit, the row. Good schemas start from a written list of reads, turn each read into a contiguous key range or a single row, encode keys as short fixed-width binary components in unsigned sort order, use reverse timestamps for newest-first reads, prefer tall rows for anything that grows and wide rows for bounded data that must change together, and treat every cross-row or cross-table write as non-atomic, with ordered, idempotent writes and a repair job. Hotspotting, column families and TTL are separate decisions made on top of that foundation.