HBase: Strong Consistency on Hadoop
Your team already runs Hadoop — HDFS holds the lake, Spark crunches it nightly (posts 2–5), and Cassandra from the last post absorbs the write firehose. Then product drops the requirement that breaks the plan: the user-profile service needs millisecond point reads on two billion profiles, and every read must reflect the latest write. Not eventually. Now. Cassandra can be tuned toward strong consistency, but the team asks a sharper question: is there a store that gives strong row-level consistency by construction, on the HDFS cluster we already operate? That's HBase — Google's BigTable design, rebuilt on Hadoop. And it comes with a catch that defines the entire post: in HBase, the row key is the schema, and a bad one doesn't just slow you down — it funnels your whole cluster's writes through a single machine.
Lab honesty, up front. The labs in this post are marked [NOT RUN]: this sandbox blocks raw TCP from Java processes, so HBase's JVM daemons can't communicate here; run this on any normal machine. (Verified with a socket probe: a Java process connecting to loopback TCP receives a policy message instead of reaching the server, while the same connection from Python works. Three start-hbase.sh attempts behaved exactly that way — daemons start, inter-process RPC never connects.) What did run here: the Java client program below compiles cleanly against the real hbase-client 2.5.16 jars (verified with javac), and the two storage defaults quoted below were read from this release's own hbase-default.xml. No shell output, region assignment, or client result below is fabricated — the lab shows commands, not invented results.
What HBase is
HBase is the open-source implementation of Google's BigTable paper: a wide-column store that lives on HDFS. A table is a sparse, sorted map — rows identified by a row key (an arbitrary byte array), columns grouped into column families (fixed at table creation), with dynamic column qualifiers inside each family, and every cell carrying a timestamp so multiple versions of a value coexist. Absent columns cost nothing — a row with three populated qualifiers stores three cells, not a full schema's worth of nulls.
Compared with Cassandra (last post): Cassandra is masterless — any node can coordinate a request, and you tune consistency per operation (ONE, QUORUM, ALL). HBase takes the opposite architectural bet: exactly one region server owns each slice of the key space, so every operation on a row goes to that row's owner. Strong per-row consistency falls out of the architecture instead of being a tuning decision. You'll sometimes hear "Cassandra leans AP, HBase leans CP" — treat that strictly as architectural shorthand, not the consistency model. The model is: single owner per row range.
Architecture: tables, regions, region servers
Rows are sorted by row key as raw bytes. A table is carved into regions — contiguous ranges of that sorted key space — and each region is served by exactly one region server. When a region grows past a configured size threshold (10 GB by default in current releases — read from this release's hbase-default.xml, not from memory), it splits in two. Three more pieces complete the picture:
- HMaster — assigns regions to region servers, reassigns them on failure, handles table creation and schema changes. It is not in the read/write path: clients talk to region servers directly.
- ZooKeeper — holds the active master's address, tracks region-server liveness through ephemeral znodes, and anchors the
/hbaseznode tree. When a region server dies, its ephemeral node disappears and the master reassigns its regions. - hbase:meta — a regular HBase table that acts as the catalog: which region covers which key range, and which region server owns it. Clients look this up once, cache it, and then go straight to the owning region server.
Notice the failure semantics this creates: if Region Server 1 dies, region A's rows are unavailable until the master notices (via ZooKeeper) and reassigns the region elsewhere — the data itself is safe on HDFS, but reads for that range stall. That is the price of single ownership, and it's exactly what the "leans CP" shorthand is gesturing at: under a partition, HBase prefers to withhold a row range rather than serve a possibly stale copy from somewhere else.
The row key is the schema
Principle: the row key is the schema. In a relational database the schema describes the data and the query planner figures out access. In HBase there is no query planner and, out of the box, no secondary index: the efficient operations are a point Get on a full row key and a Scan over a contiguous row-key range. Everything else is a full-table scan. So the access pattern must be designed into the key itself — get the key wrong and no amount of tuning saves you.
Row keys sort as unsigned bytes, lexicographically. Regions are contiguous slices of that ordering. Now watch what two designs do to a cluster ingesting events:
Three techniques, each with an honest cost:
- Salting (spread the writes). Prefix the key with
hash(entity) mod N, e.g.a3#20261006T120000#user123. Writes scatter across N buckets instead of piling onto one region. The cost: a read for one user now needs N range scans (one per bucket) unless the salt is computable from the query — which it is here, since it's derived from the user id. - Reversing the timestamp (newest first). Use
Long.MAX_VALUE - timestampso a scan from the key prefix returns the newest events first — the common "last 50 events for user X" query. Be precise about what this does not do: it does not fix hotspotting. Reversed timestamps are still monotonic, so writes still pile onto one region — now at the low end instead of the high end. Combine reversing with salting when you need both. - Entity-first ordering (keep related rows adjacent).
user123#20261006T120000puts one user's events in neighboring rows: a single range scan reads the whole history. This distributes well across regions as a rule of thumb when the leading ids are well distributed (hashes, UUIDs); it hotspots only if the leading component clusters (e.g. one giant tenant's id on every key — then salt that).
Two more schema facts that shape key design. First, column families are fixed at table creation but qualifiers are dynamic — you declare e once, then write e:login, e:purchase, e:anything-new freely. Second, pre-split the table instead of starting with one region and waiting for splits: create 'user_events', 'e', SPLITS => ['4', '8', 'c'] (hex-string split points for salted keys) starts you with four regions on day one. An empty single-region table that must split its way up under load is a self-inflicted hotspot.
(Apache Phoenix adds a SQL layer over HBase, but it doesn't change the physics: its queries still resolve to row-key gets and range scans underneath. The key design rules above apply either way.)
Strong per-row consistency, by construction
Here's the mechanism the shorthand points at. Because one region server owns each region, every operation on a given row is handled by that row's owner — there is no second copy that could answer differently. A Put to a row is atomic across all its column families, and a Get sees the row as of a single consistent point. Read-your-write holds without any tuning: you write to the owner, you read from the owner.
ZooKeeper's role in this is coordination, not data: it elects the active master, watches region-server liveness via ephemeral znodes, and anchors the metadata so a failed server's regions get reassigned to a live one. The actual consistency guarantee comes from single ownership, with ZooKeeper making failover fast enough to matter.
The honest boundary: there are no multi-row transactions. Two rows may live in different regions on different servers, and HBase will not atomically update both. If your write pattern needs "debit A and credit B or neither," that's a relational database's job (or application-level idempotency, the way the streaming post handles exactly-once sinks). HBase gives you atomic, strongly consistent single rows at massive scale — nothing more, and that's usually enough for event and profile workloads.
Stack it against Cassandra from the last post and the choice sharpens: Cassandra's masterless design lets any node coordinate and leaves consistency to you per operation — flexible, and it leans toward staying available when parts of the cluster can't talk to each other. HBase's single-owner design gives you strong row consistency without thinking about it, and charges you with unavailability of a row range while its owner is being replaced. Neither is "more consistent" in the abstract; they place the same tradeoff in different spots.
LSM-tree storage: memstore, HFiles, compaction
HBase never updates data in place — it shares the log-structured merge-tree design with Cassandra (last post), so this will feel familiar. The write path, briefly and precisely:
- Write-ahead log (WAL). Every mutation is appended to the region server's WAL first. This is the durability story: if the server dies before data reaches a file, the WAL is replayed.
- Memstore. The mutation also lands in the region's memstore — an in-memory sorted buffer, one per column family per region. Reads check here first, so recent writes are fast.
- Flush. When a memstore passes its size threshold (128 MB by default — again, read from this release's
hbase-default.xml), it's flushed as an immutable HFile, HBase's sorted file format, written to HDFS. - Compaction. HFiles accumulate, and a read would have to merge dozens of them. Minor compaction merges a few small HFiles into a larger one; major compaction merges all of a region's HFiles and, while doing so, drops cells that were deleted (tombstones) or expired by TTL. Major compactions are I/O-heavy, which is why they're typically scheduled, not left to chance.
Reads merge the memstore, the relevant HFiles, and the block cache into one consistent view of the row. Deletes are just another write — a tombstone marker that the next major compaction garbage-collects. The whole design optimizes for the workload HBase serves: heavy random writes, fast point reads, and scans that stream sorted files sequentially.
HBase on HDFS: the honest layering
Everything durable in HBase — HFiles, the WAL — lives under hbase.rootdir, which on a real cluster is an HDFS path. That layering is the point: HBase doesn't reimplement replication or failure recovery for storage; it inherits them from HDFS (post 2). A region server dies, its WAL is sitting on HDFS, the master reassigns its regions, and the new owner replays the log. Compute (region servers) is disposable; storage (HDFS) is durable.
Standalone-mode honesty: the lab below sets hbase.rootdir to the local filesystem, exactly as HBase's own standalone docs prescribe. Same code paths, zero distribution — one JVM plays master, region server, and ZooKeeper. It teaches the API and the data model faithfully; it teaches you nothing about region distribution or failover. Don't mistake "it ran on my laptop" for "it works at scale" — that confusion is how bad row keys reach production.
Lab [NOT RUN]: standalone HBase, shell, and the Java client
Run this on any normal machine with JDK 8, 11, or 17. (HBase 2.5 officially targets 8/11; on 17 it starts with a set of --add-opens flags in HBASE_OPTS covering java.base/java.nio, java.base/java.lang, java.base/java.util and friends — the author's attempts here got past module access cleanly, so this is configuration, not a blocker.)
Install and start. Download 2.5.16 from the Apache archive (archive.apache.org/dist/hbase/2.5.16/), extract it, and point hbase-site.xml at local paths for standalone mode:
$ tar xzf hbase-2.5.16-bin.tar.gz
$ cat conf/hbase-site.xml
<configuration>
<property>
<name>hbase.rootdir</name>
<value>file:///path/to/hbase-2.5.16/data</value>
</property>
<property>
<name>hbase.zookeeper.property.dataDir</name>
<value>/path/to/hbase-2.5.16/zk</value>
</property>
</configuration>
$ export JAVA_HOME=/path/to/jdk17
$ bin/start-hbase.sh
running master, logging to .../logs/hbase--master-....log
Explore in the shell [NOT RUN]. The hbase shell is a JRuby REPL — the fastest way to feel the data model. Create the table from the row-key section, write three rows with designed keys, read one back, and scan a user's range:
$ bin/hbase shell
hbase:001:0> create 'user_events', 'e'
hbase:002:0> put 'user_events', '321resu#20261006120000', 'e:login', 'web'
hbase:003:0> put 'user_events', '321resu#20261006121500', 'e:purchase', 'sku-42'
hbase:004:0> put 'user_events', '654resu#20261006120500', 'e:login', 'mobile'
hbase:005:0> get 'user_events', '321resu#20261006121500'
hbase:006:0> scan 'user_events', {STARTROW => '321resu#', LIMIT => 5}
The Java client [COMPILED, NOT RUN]. The same operations through hbase-client. The full program is HBaseLab.java alongside this draft; the core flow:
Configuration conf = HBaseConfiguration.create();
conf.set("hbase.zookeeper.quorum", "localhost"); // standalone lab
try (Connection conn = ConnectionFactory.createConnection(conf);
Admin admin = conn.getAdmin()) {
TableDescriptor td = TableDescriptorBuilder.newBuilder(TABLE)
.setColumnFamily(ColumnFamilyDescriptorBuilder.of(CF))
.build();
admin.createTable(td);
try (Table table = conn.getTable(TABLE)) {
// Designed row keys: user id first, then the timestamp.
put(table, "321resu#20261006120000", "login", "web");
put(table, "321resu#20261006121500", "purchase", "sku-42");
// Point read: one row, read-your-write.
Result r = table.get(new Get(Bytes.toBytes("321resu#20261006121500")));
// Range scan: one user's events, in row-key order.
Scan scan = new Scan()
.withStartRow(Bytes.toBytes("321resu#"))
.withStopRow(Bytes.toBytes("321resu$"));
try (ResultScanner scanner = table.getScanner(scan)) {
for (Result row : scanner) { /* ... */ }
}
}
}
What was verified here: javac -cp "$HBASE_HOME/lib/*" HBaseLab.java compiles cleanly against hbase-client 2.5.16 — every API call above ( ConnectionFactory, TableDescriptorBuilder, Put/Get/Scan, withStartRow/withStopRow) type-checks against the real jars. Execution is [NOT RUN] for the sandbox reason above — run it where TCP works and the shell transcript will match these commands.
The hotspot experiment to try on a real machine. Write one million rows with timestamp-first keys, then watch the region server metrics: one server's request count dwarfs the rest. Repeat with a salted prefix and watch the load even out — then run a single-user read and feel the fan-out cost. That experiment, more than any paragraph here, is what "the row key is the schema" means.
When HBase, when Cassandra: qualified guidance
Both are wide-column stores; the last two posts gave you both architectures. The decision is rarely about features — it's about what you already run and which guarantee you want by default:
- Rule of thumb, not a law: already on Hadoop + need strong row consistency → HBase is the natural fit. If HDFS, YARN, and Spark are already your platform and the workload is versioned time-series or profile data with key-scoped reads, HBase slots in with no new storage layer to operate.
- Rule of thumb, not a law: need active-active writes across data centers, or no Hadoop footprint → Cassandra tends to fit better. Cassandra's masterless design and tunable per-operation consistency were built for that shape; dragging a Hadoop cluster along just for storage is the wrong trade.
- Neither, when the query pattern says so. Ad-hoc analytics over billions of rows is ClickHouse's question, sub-millisecond serving is Redis's — that's the next post. "HBase or Cassandra?" with no query pattern attached is unanswerable; write the query first, as the first post's principle says.
And the operational footnote both choices share: wide-column stores move schema design from DDL into the key. Teams that treat HBase "like Postgres with one big table" get hotspots and full scans; teams that design the row key from the query list get a store that quietly does its job for years.
Cheat sheet
- HBase is BigTable on HDFS: sparse sorted map, row key → column family → qualifier → timestamped cells.
- Tables split into regions (contiguous key ranges, 10 GB split threshold by default); exactly one region server owns each region.
- HMaster assigns regions and handles DDL — never in the data path. ZooKeeper handles master election and server liveness; hbase:meta maps regions to owners.
- The row key is the schema. Point gets and range scans are the efficient ops; design the key from the query pattern.
- Timestamp-first keys hotspot (monotonic → one region takes all writes). Salt with
hash(entity) mod Nto spread writes; reverse timestamps for newest-first scans (doesn't fix hotspots); entity-first keys keep related rows adjacent. - Single-row operations are atomic with read-your-write, by construction (single owner per row). No multi-row transactions — that's the boundary.
- "Leans CP" is architectural shorthand: a dead region server's range is unavailable until reassigned. Cassandra (last post) leans the other way with tunable per-operation consistency.
- LSM storage: WAL → memstore (flush at 128 MB) → immutable HFiles on HDFS; minor/major compaction merges files and collects tombstones.
- Durability and replication come from HDFS, not HBase. Standalone mode is local-FS only — same code paths, zero distribution.
- Java client:
ConnectionFactory→Admin/Table→Put/Get/Scanwith byte-array keys. Compiles against hbase-client 2.5.16; needs JDK 8/11 (or 17 with--add-opens).
What's next
You now hold both wide-column answers: Cassandra's tunable, masterless write machine and HBase's strongly consistent row store on Hadoop. But downstream serving has two more shapes this track promised — and they're nothing like these. Redis answers "what's in the cache right now?" in microseconds from memory, and ClickHouse answers "how did revenue trend this quarter?" by reading columns instead of rows.
The next post, Redis & ClickHouse: Storage & Serving, puts both on the bench: Redis data structures as a serving layer in Java, ClickHouse's MergeTree engine chewing through billions of rows, and the decision rule for which downstream question each one owns.
Field check before you move on: (1) Design row keys for an order-events table queried two ways — "all events for order X" and "latest 100 events across all orders" — and explain which query your key favors and what the other one costs. (2) Explain, in your own words, why reversing a timestamp gives newest-first scans but does not fix a write hotspot. (3) On a real machine, run the shell lab above, then pre-split a copy of the table and compare region counts via the HBase web UI. (4) Write one paragraph recommending HBase or Cassandra for a fraud-signal store that must reflect each write immediately and already runs on the company's Hadoop cluster — then argue the opposite choice and name the assumption you'd have to change.
Continue: Java Learning Roadmap 2026
Comments
Post a Comment