HBase is not a graph database, yet a lot of large graphs live in it: social follow graphs, device and identity graphs for fraud detection, web link graphs, and knowledge graphs built by batch pipelines. The reason is scale. HBase stores billions of rows across hundreds of servers, sustains heavy write rates, and sits next to Hadoop and Spark so the same data feeds both online lookups and offline analytics. What it lacks is any notion of a vertex, an edge or a traversal. You build those out of sorted rows, column families and batched reads.

This article shows how to do that well. It covers the two basic layouts, row key design, supernodes, how a multi-hop query executes, why writing an edge is not atomic, and when to use JanusGraph instead of rolling your own. A worked example sizes a friends-of-friends query, and the article ends with a checklist.

What HBase gives a graph

Start from what HBase gives you. A table is a sorted map from row key to a set of column families; each family holds any number of columns, named by arbitrary byte qualifiers, each holding versioned cells. Rows are sorted by key and split into regions served by region servers. Two properties matter most for graphs. First, a single row is the unit of atomicity: all mutations to one row in one call succeed or fail together, and nothing spans rows. Second, reads are cheap when they are contiguous: a Get of one row or a Scan over adjacent keys touches one region, while scattered keys mean scattered RPCs.

A graph query is mostly one operation repeated: given a vertex, list its neighbours, optionally filtered by edge label and direction. The whole design goal is to make that operation a single-row read, or at worst a short scan. The general rules for keys, families and row width are covered in HBase schema design and wide versus tall tables; here they are applied to vertices and edges.

Two layouts: adjacency rows and edge tables

Layout A, the adjacency row. Each vertex is one row. A small family, call it p, holds vertex properties. A second family, e, holds one column per incident edge. The qualifier encodes direction, label and the other endpoint, so that edges sort into useful groups: all outgoing follows, then all outgoing likes, then all incoming follows, and so on. The cell value holds the edge properties, or is empty if there are none.

Adjacency lists on HBase: one row per vertex, edges as columns, both directions storedRow: salt + vertex ide.g. 3f|user:42family p: propertiesname, created, ...family e: O|follows|user:77value: edge propsfamily e: I|follows|user:9value: edge propsSupernode rowsuser:1#b00 ... user:1#b63high-degree vertex split into bucketsTraversal: frontier of vertex ids, one batched multi-get per hopHop 01 rowHop 1~200 rowsHop 2~40,000 rowsHop 3cap or sampleGet listGet listlimitEach hop multiplies reads by the average degree; budget fan-out before you query
Layout A: a vertex row with a property family and an edge family whose qualifiers sort by direction and label. High-degree vertices are split into bucket rows. A traversal turns each hop into one batched multi-get.

Because qualifiers within a row are sorted, the neighbours of user 42 via outgoing follows edges are a contiguous range of columns, which a ColumnPrefixFilter on O|follows| returns without reading the rest of the row. Every edge is stored twice, once as an out-edge in the source row and once as an in-edge in the destination row. That doubles storage and write volume, but it makes both directions a single-row read, and it is what lets you delete a vertex cleanly, because the vertex row lists every other row that references it.

Layout B, the edge table. Each edge is its own row with key src|label|dst. All out-edges of a vertex are adjacent rows, so neighbours are a prefix Scan. In-edges need a second table keyed dst|label|src. This layout never produces a giant row, and paging through a high-degree vertex is natural, but a neighbour lookup is a Scan rather than a Get, which costs more per call, and vertex properties live in yet another table.

ConcernA: vertex rowB: edge table
Neighbour lookupone Get, filtered by qualifier prefixprefix Scan
High-degree verticesrow grows without bound; needs bucketingnaturally spread across rows
Vertex plus edges togethersame row, one readtwo tables, two reads
Atomicityvertex and its local edge copies in one roweach edge row alone
Typical fitmost vertices have modest degreedegree is heavy-tailed and unbounded

Many production systems use A for ordinary vertices and B, or bucketed A, for the small fraction of vertices whose degree explodes.

Row keys and qualifier encoding

Vertex ids from upstream systems are often sequential or share prefixes, which sends all new vertices to the same region; see hotspotting. Prefix the key with a short hash of the id, for example the first byte of an MD5, so writes spread across regions while each vertex is still found by recomputing the hash. Keep the label and direction encoding compact, because every byte of the qualifier is repeated in every cell on disk: map labels to a one- or two-byte code through a small dictionary, and encode ids as fixed-width binary rather than strings when you can.

// Writing one edge (src -follows-> dst) in layout A with the HBase Java client.
static final byte[] E = Bytes.toBytes("e");

static byte[] rowKey(String vertexId) throws NoSuchAlgorithmException {
    byte[] id = Bytes.toBytes(vertexId);
    byte salt = MessageDigest.getInstance("MD5").digest(id)[0];
    return Bytes.add(new byte[] { salt }, id);
}

static byte[] qualifier(char dir, short label, String other) {
    return Bytes.add(new byte[] { (byte) dir }, Bytes.toBytes(label), Bytes.toBytes(other));
}

void addEdge(Table t, String src, short label, String dst, byte[] props, long ts) throws Exception {
    Put out = new Put(rowKey(src)).addColumn(E, qualifier('O', label, dst), ts, props);
    Put in  = new Put(rowKey(dst)).addColumn(E, qualifier('I', label, src), ts, props);
    // Two rows, two independent mutations: HBase gives no cross-row atomicity.
    t.batch(List.of(out, in), new Object[2]);
}

Passing an explicit timestamp matters: if a retry rewrites the same edge, both copies carry the same version, so the write is idempotent rather than producing a second version.

Supernodes and bucketing

Real graphs have heavy-tailed degree distributions. A celebrity account, a shared corporate Wi-Fi address in an identity graph, or a popular product can have millions of edges. In layout A that is one row with millions of columns, and a row cannot be split across regions. The consequences are concrete: a full Get of that row can exhaust client or server memory, compactions rewrite the whole row, and every query that touches the vertex lands on one region server. Wide row limits covers the mechanics.

The fix is bucketing. Once a vertex's degree passes a threshold, store its edges in rows keyed salt|vertexId#bNN, where NN is the hash of the other endpoint modulo a bucket count such as 64. The base row keeps properties and a small marker column recording the bucket count. A reader that sees the marker issues a Get per bucket in one batch; a writer computes the bucket from the destination id. Because bucket rows have different keys they hash to different regions, which spreads the load as well as the size.

Two more defences belong in every reader. Never fetch an unbounded row: use Scan.setBatch or ColumnPaginationFilter to page through columns, and set setAllowPartialResults(true) on scans of wide rows. And cap fan-out at the application level: a traversal that reaches a supernode usually wants a sample or the top few edges by weight, not all of them. A per-vertex degree counter, kept with an HBase Increment, lets queries decide before reading.

Running a traversal

HBase has no query planner for graphs, so the client runs the traversal. The efficient pattern is breadth-first with batched reads: keep a frontier of vertex ids, turn the whole frontier into a list of Gets restricted to the right qualifier prefix, send them in one Table.get(List<Get>) call, which the client groups by region server, and build the next frontier from the results.

def k_hop(table, start, label_code, hops, max_frontier=50_000, max_per_vertex=500):
    """Breadth-first k-hop neighbourhood over layout A. Returns {vertex: depth}."""
    seen = {start: 0}
    frontier = [start]
    for depth in range(1, hops + 1):
        gets = [get_out_edges(v, label_code, limit=max_per_vertex) for v in frontier]
        nxt = []
        for batch in chunks(gets, 1000):            # bound each RPC fan-out
            for v, neighbours in table.multi_get(batch):
                for n in neighbours:
                    if n not in seen:
                        seen[n] = depth
                        nxt.append(n)
        if len(nxt) > max_frontier:                  # stop the explosion explicitly
            nxt = sample(nxt, max_frontier)
        frontier = nxt
        if not frontier:
            break
    return seen

Here get_out_edges builds a Get with a ColumnPrefixFilter for the direction and label and a ColumnPaginationFilter for the limit, and expands bucketed vertices into several Gets. The two caps are not optional: without them, one query that wanders into a dense region can issue millions of reads and starve every other tenant.

Worked example: friends of friends at 50 million users

Take a follow graph with 50 million users and an average of 200 follows each. Out-edges alone are 10 billion columns; storing both directions makes 20 billion cells. With a 1-byte salt and an 8-byte binary user id the row key is 9 bytes, and a qualifier made of a 1-byte direction, a 2-byte label code and the 8-byte neighbour id is 11 bytes, and an empty value keeps cells small, though HBase's per-cell key overhead (row, family, timestamp, type) still dominates, so compression and data block encoding on the e family are worth enabling.

Now the query: suggest accounts followed by the people a user follows. Hop 1 is one Get returning about 200 ids. Hop 2 is 200 Gets, batched into a single multi-get call that the client splits by region server, returning about 40,000 edges. That is fine for an online request if each Get is limited and cached blocks are warm. A third hop would mean around 8 million reads, which is not an online query on any store; precompute it in batch with Spark, for example with GraphFrames, and write the results back to HBase as a materialised recommendation row. If one of the 200 accounts is a celebrity bucketed into 64 rows, cap that vertex at the top few hundred edges so it does not dominate the result.

Consistency: the half-written edge

Because the out-edge and the in-edge live in different rows, a crash between the two mutations leaves a half-written edge: the source thinks it follows the destination, but the destination does not list the follower. HBase core has no multi-row transactions; the MultiRowMutationEndpoint coprocessor gives atomicity only for rows in the same region, which two random vertices almost never are. Accept that and design for repair.

  • Write in a fixed order and make retries idempotent. Write the copy that drives user-visible reads first, use the same explicit timestamp for both copies, and retry the whole pair on failure.
  • Log intent. Append each edge mutation to a durable log such as a Kafka topic before applying it, and have the writer replay unacknowledged entries after a crash.
  • Reconcile in batch. A periodic Spark job scans the e family, emits every out-edge and in-edge as a pair, and fixes any edge present in only one direction. Track the count of repaired edges as a health metric; a rising count means a writer bug.
  • Deletes are tombstones. Removing a vertex means deleting its row and every copy of its edges in neighbours' rows, using the in- and out-edge columns as the list. Deleted cells occupy space until a major compaction removes them.

Reads can also see one copy without the other briefly even when nothing crashes, so code that requires symmetric edges should check, not assume.

JanusGraph on HBase

If you need a real graph query language, JanusGraph is an open-source graph database that implements Apache TinkerPop and can use HBase as its storage backend. It stores each vertex's adjacency list as a row in an edge store, sorted so that label- and direction-restricted reads are range reads, and keeps composite indexes in a separate index store, which is the same idea as layout A with the encoding done for you. Configuration is a properties file:

storage.backend=hbase
storage.hostname=zk1.example.internal,zk2.example.internal,zk3.example.internal
storage.hbase.table=janusgraph
# optional: full-text and range queries go to an external index
index.search.backend=elasticsearch
index.search.hostname=es1.example.internal

You then write Gremlin, for example g.V().has('user','uid',42).out('follows').out('follows').dedup().limit(100), and JanusGraph turns each step into reads against HBase. The same physics apply: multi-hop traversals multiply reads and supernodes hurt, though JanusGraph's vertex-centric indexes help with the latter. Check the JanusGraph release notes for the HBase versions a given release supports before deploying.

Failure modes

  • Unbounded row reads. A Get on a supernode returns millions of cells and times out or exhausts heap. Always page and cap.
  • Region hotspots. Unsalted sequential ids, or one viral vertex, put all traffic on one region server. Salt keys and bucket high-degree vertices.
  • Asymmetric edges. Crashes or partial batch failures leave one copy. Run reconciliation and alert on its repair count.
  • Frontier explosion. A three-hop query on a dense graph issues millions of Gets. Set hard frontier and per-vertex caps, and move deep traversals to batch.
  • Tombstone build-up. Graphs with heavy edge churn accumulate deletes that slow reads until a major compaction. Watch read latency on high-churn families and schedule compactions.

What to do next

  1. Write down the three queries that matter, with the maximum hops and acceptable latency for each.
  2. Measure the degree distribution of your graph and pick a supernode threshold from its tail.
  3. Implement layout A with salted keys, compact qualifiers and explicit timestamps; add bucketing above the threshold.
  4. Build the traversal with batched multi-gets, a per-vertex limit and a frontier cap, and load-test it against a known supernode.
  5. Add an intent log and a nightly reconciliation job, and chart the number of repaired edges.
  6. Precompute anything beyond two hops in Spark and serve it from a materialised row.
  7. If you need ad hoc Gremlin, prototype JanusGraph on the same HBase cluster and compare its read counts with your hand-built layout.
Key takeaway: HBase can hold very large graphs if you shape the data around single-row neighbour lookups. Store each vertex as a row with direction- and label-prefixed edge columns, salt the keys, bucket high-degree vertices, and traverse breadth-first with batched, capped multi-gets. Edges written to two rows are not atomic, so make writes idempotent and reconcile in batch. Push anything deeper than two hops to Spark, and reach for JanusGraph when you need Gremlin rather than a fixed access pattern.