A session store looks like the simplest workload in the world: put a blob under a key, get it back on the next request, throw it away when the user leaves. Most teams reach for Redis or Memcached, and for many sites that is right. But some session workloads outgrow a cache. Sessions that must survive a data centre restart, conversation state for an assistant that users resume days later, device sessions for tens of millions of users that security teams want to revoke in bulk, and click-stream sessions that analysts want to scan with Spark all push toward a durable, horizontally scaled store. If you already run HBase, it can do this job well, provided you design around how it actually stores, expires and serves data.
This article builds an HBase session store from first principles: the row key, the column family layout, how expiry really works, concurrent updates, the per-user index needed for logging out everywhere, and the operational behaviour you will see when a RegionServer dies. It ends with a sizing worked example and a checklist.
Architecture at a glance
When HBase is the right session store
Pick the store by what the session has to survive and who else needs to read it. A cache-only store loses sessions on eviction or restart, which is fine when re-login is cheap. HBase earns its place when you need three things at once: durability through the write-ahead log, linear scale by adding RegionServers, and cheap bulk access for analytics or revocation. It also gives you strongly consistent single-row reads and writes, because each row is served by exactly one region at a time. That matters for sessions: a read after a write from the same user always sees the write, with no eventual-consistency window to reason about.
What HBase does not give you is cache-like tail latency under all conditions. Reads that miss the block cache touch HFiles on HDFS, Java garbage collection adds pauses, and a RegionServer failure makes its regions unavailable until they are reassigned. Design for single-digit-millisecond medians and plan explicitly for the tail.
Row key design
The row key decides both distribution and security. Session tokens are usually 128 or more bits of randomness, so they are already uniformly distributed and need no salting: rows spread evenly across pre-split regions. Do not, however, use the raw token as the key. The token is a bearer credential, and anyone holding a table dump, a snapshot export or debug logs of row keys could replay it. Store SHA-256 of the token instead. Lookups hash the incoming cookie and do a point Get, and a leaked table reveals no usable tokens.
Avoid time-ordered identifiers such as UUIDv7 or snowflake IDs as keys. Their prefixes increase monotonically, so every new session lands in the last region and one RegionServer takes all writes, the classic hotspot. Hashing removes that, at the cost of losing range scans by creation time, which a session table does not need. Because hashed keys are uniform, pre-splitting is easy: divide the first byte or two into equal ranges and create the table with those split points so load is spread from the first request.
Column family layout
Keep one small column family, here called s, because every family is a separate store with its own memstore and files, and a session row is read whole. Inside it, store the session state as one serialized cell (JSON, Protobuf or Avro), plus a ver cell used for concurrency control. Set VERSIONS => 1 since old states have no value, enable a ROW bloom filter so Gets skip HFiles that cannot contain the row, and consider IN_MEMORY if the table is small relative to the block cache. Compress with a fast codec; session JSON compresses well.
TableDescriptor desc = TableDescriptorBuilder.newBuilder(TableName.valueOf("sessions"))
.setColumnFamily(ColumnFamilyDescriptorBuilder.newBuilder(Bytes.toBytes("s"))
.setMaxVersions(1)
.setTimeToLive(8 * 60 * 60) // seconds: absolute ceiling, the garbage collector
.setBloomFilterType(BloomType.ROW)
.setCompressionType(Compression.Algorithm.SNAPPY)
.setDataBlockEncoding(DataBlockEncoding.FAST_DIFF)
.build())
.build();
// 16 equal regions on the first byte of sha256(token)
byte[][] splits = new byte[15][];
for (int i = 1; i < 16; i++) {
splits[i - 1] = new byte[] { (byte) (i * 16) };
}
admin.createTable(desc, splits);
How expiry really works
Expiry is where most HBase session stores go wrong. TTL in HBase is evaluated per cell, against the cell timestamp. A family TTL of 1,800 seconds means each cell becomes invisible 30 minutes after it was written. Readers stop seeing expired cells immediately, but the bytes stay on disk until a compaction rewrites the files, and a store file whose cells have all expired can be dropped whole.
Now consider a sliding 30-minute idle timeout implemented with a 1,800-second family TTL and several cells per session: user id, cart, preferences, last_seen. Each request updates only last_seen. Thirty minutes after login, the user id and cart cells expire even though the user is active, and the session silently loses half its state. That is not an HBase bug; it is per-cell semantics. There are two correct designs.
First, the single-cell design used here: every update rewrites the whole state cell, so its timestamp always reflects the last activity and the TTL behaves like a sliding timeout. Second, the application-expiry design: the family TTL is set to the absolute maximum session lifetime (say eight hours) and acts only as a garbage collector, while the application compares last_seen with the idle limit on every read and treats stale sessions as missing. Most production systems combine both: one rewritten cell, an idle check in code, and a generous family TTL as a safety net.
Per-cell TTL is also available through Mutation.setTTL, which takes milliseconds while the family setting takes seconds; mixing them up gives sessions that last a thousand times too long or vanish instantly. A cell TTL can shorten a lifetime but cannot extend it beyond the family TTL. Prefer explicit expiry over Delete for logout too: a Delete writes a tombstone that every read must skip until a major compaction removes it, so heavy delete traffic slowly increases read cost.
Concurrent updates
Two browser tabs or two app instances can update one session concurrently. A blind Put means last writer wins, which loses a cart item. HBase offers atomic single-row operations: Increment, Append and conditional mutation. The pattern is optimistic concurrency: read state and version, modify in memory, then write only if the version is unchanged. HBase 2.4 introduced the CheckAndMutate builder; the older fluent table.checkAndMutate(row, family) form is deprecated.
public boolean updateSession(Table t, byte[] key, Function<State, State> change)
throws IOException, InterruptedException {
for (int attempt = 0; attempt < 5; attempt++) {
Result r = t.get(new Get(key).addFamily(S));
if (r.isEmpty()) return false; // expired or revoked
long ver = Bytes.toLong(r.getValue(S, VER));
State next = change.apply(State.decode(r.getValue(S, STATE)));
if (next.idleSeconds(clock) > IDLE_LIMIT) return false;
Put put = new Put(key)
.addColumn(S, STATE, next.encode()) // rewrite whole state: TTL slides
.addColumn(S, VER, Bytes.toBytes(ver + 1));
put.setDurability(Durability.SYNC_WAL);
CheckAndMutate cam = CheckAndMutate.newBuilder(key)
.ifEquals(S, VER, Bytes.toBytes(ver))
.build(put);
if (t.checkAndMutate(cam).isSuccess()) return true;
Thread.sleep(ThreadLocalRandom.current().nextInt(2, 10 << attempt)); // jittered backoff
}
throw new ConflictException("session busy");
}Rewriting both cells on every successful update keeps their timestamps aligned, so neither expires before the other. For hot counters such as request counts, use Increment instead, which is atomic without a read.
Durability per mutation
Each mutation can choose a durability level. SYNC_WAL waits for the WAL edit to be written to the HDFS pipeline before acknowledging; ASYNC_WAL acknowledges first and syncs shortly after, so a RegionServer crash can lose the last moments of writes; SKIP_WAL writes only to the memstore and loses everything since the last flush on a crash. For sessions, use SYNC_WAL for login, privilege changes and revocation, because losing a revocation re-enables a session someone deliberately killed. ASYNC_WAL is a reasonable trade for last_seen-style touches where losing a few seconds just means an idle timer is slightly early. Never use SKIP_WAL for session state.
Log out everywhere: the user index
Security teams eventually ask for log out everywhere: revoke every session of a user after a password change. With keys that are token hashes, the sessions table cannot answer which sessions belong to user 42. Add a second table, user_sessions, keyed by a two-byte salt of the user id, then the user id, then reversed creation time, with the session key as the qualifier. A prefix scan returns the user's sessions newest first.
HBase has no multi-row transactions across tables, so order the writes so failures are harmless. On login, write the index entry first and the session second: a crash between them leaves an index entry pointing at nothing, which the revocation job simply skips. On revocation, delete or overwrite the sessions first, then clean the index. Give the index family the same TTL ceiling so orphans age out. Do not try to keep the index perfectly in sync; make every reader tolerate dangling references.
def logout_everywhere(conn, user_id):
prefix = salt(user_id) + user_id.encode() + b"|"
idx = conn.table("user_sessions")
sess = conn.table("sessions")
keys = [q for _, data in idx.scan(row_prefix=prefix) for q in data] # qualifiers = session keys
with sess.batch(batch_size=500) as b: # session rows first: revocation is the security event
for k in keys:
b.delete(k.split(b":", 1)[1])
with idx.batch(batch_size=500) as b:
for row, _ in idx.scan(row_prefix=prefix):
b.delete(row)The snippet uses a HappyBase-style Thrift client; in Java the same flow is a prefix Scan plus a batched list of Deletes. These deletes create tombstones, which is acceptable because revocation is rare compared with ordinary expiry.
Worked example: sizing
Size a store for 20 million sessions active at peak, a 2 KB serialized state, an idle timeout of 30 minutes, an absolute limit of 8 hours, and an average of one state update every 20 seconds per active session.
Writes: 20,000,000 / 20 = 1,000,000 updates per second at peak, each rewriting roughly 2 KB, so about 2 GB per second of logical writes into memstores and WALs, before HDFS replication triples the WAL bytes on disk. That is a big cluster. The first design lever is therefore not HBase tuning but write reduction: debounce touches so last_seen is written at most once a minute unless state changed, and the rate typically drops by an order of magnitude. Reads: one Get per request, served mostly from the block cache because recently written sessions are hot.
Storage: live data is 20,000,000 x 2 KB = 40 GB, but with one rewrite per 20 seconds, each session produces many superseded versions between compactions. VERSIONS => 1 hides them, yet they occupy store files until compaction. Expect several times the live size on disk in steady state, and watch compaction throughput: if compactions fall behind, store file counts climb, reads slow, and eventually writes block. A per-region write rate figure from your own load test, not a rule of thumb, should decide the region count.
Failure modes
The failure modes are predictable, which means you can rehearse them.
- RegionServer crash. Its regions are offline until the master detects the failure, splits the WAL and reopens the regions elsewhere. Sessions on those regions fail lookups for that window. Region replicas let clients read a secondary copy with timeline consistency; for sessions that is usually acceptable for reads, while writes still wait for the primary.
- Silent state loss from per-cell TTL. Described above; caught by a test that keeps a session active past the TTL and asserts every field survives.
- Hotspots from sequential keys. One region's request count far above the rest. Fix the key, then split or rebuild.
- Tombstone build-up. Logout-heavy designs that Delete instead of letting TTL expire slow scans until major compaction.
- Compaction debt. Update-heavy sessions create many small files; if the store file count reaches the blocking threshold, flushes stall and writes block.
- GC pauses. Long pauses can exceed the ZooKeeper session timeout and get the RegionServer declared dead. Use off-heap BucketCache and a modern collector.
Trade-offs
| Choice | Gain | Cost |
|---|---|---|
| Single state cell | TTL slides correctly, one read | Whole blob rewritten per update |
| Hashed token key | Even load, no token leak | No range scans by time |
| SYNC_WAL for all | No lost revocations | Higher write latency |
| Region replicas | Reads during failover | More memory and storage |
| Delete on logout | Immediate removal | Tombstones until major compaction |
The headline trade-off is against a cache. Redis gives lower tail latency and simpler operations; HBase gives durability, scale and analytics on the same data. A common hybrid puts a short-TTL cache in front of HBase, invalidated on write, with HBase as the source of truth.
What to do next
- Decide what a session must survive; if the answer is nothing, use a cache instead.
- Key the table by SHA-256 of the token and pre-split on its first byte.
- Store state in one rewritten cell with VERSIONS => 1 and a ROW bloom filter.
- Set the family TTL to the absolute lifetime and enforce the idle timeout in code.
- Use CheckAndMutate with a version cell for concurrent updates.
- Use SYNC_WAL for login and revocation; consider ASYNC_WAL for debounced touches.
- Add the user_sessions index with crash-safe write ordering.
- Load test, then kill a RegionServer and measure lookup failures during recovery.
- Read more on TTL and versions, hashing row keys, WAL durability, region replicas and HBase versus DynamoDB.