Disaster recovery is the plan for the day a whole HBase cluster, or the site it runs in, stops being usable. It is different from high availability, which keeps a cluster serving when a RegionServer or a disk dies, and different from backup, which lets you go back in time after someone deletes or corrupts data. A DR plan answers two questions with numbers: how much recently written data you can afford to lose (the recovery point objective, RPO) and how long the service can be down (the recovery time objective, RTO).
This article builds a DR design from those two numbers. It maps failure classes to the HBase tools that address them, lists everything besides table data that a standby site needs, walks through bootstrapping an asynchronous standby, explains synchronous replication and its costs, and gives failover and failback runbooks with a worked example. The mechanics of the replication pipeline are covered in HBase replication, and point-in-time recovery in HBase backups.
Set RPO and RTO per table, not per cluster
A single cluster usually hosts tables of very different value. A payments ledger may need near-zero data loss, a clickstream table can lose a few minutes, and a derived feature table can be rebuilt from upstream in a day. Giving all of them the strictest target costs latency and money; giving all of them the loosest one loses data that mattered. Write the targets per table, and let them choose the mechanism.
Three tiers cover most estates. They are a starting point to argue over with the data owners, not a standard.
| Tier | RPO | RTO | Typical mechanism |
|---|---|---|---|
| Critical | Seconds or zero | Minutes | Synchronous replication, or async replication with lag alerting and a rehearsed runbook |
| Important | Minutes | Under an hour | Asynchronous replication to a warm standby |
| Rebuildable | Hours to a day | Hours | Exported snapshots plus the ability to regenerate from source |
Two numbers make the targets honest. First, measure replication lag under peak load, not on a quiet afternoon; that lag is your real RPO for async tiers. Second, time a full failover in a drill, including the human steps such as paging and decision-making; that is your real RTO.
Map failure classes to tools
Every DR tool protects against some failures and is useless for others. The most important line in this table is the last one: replication faithfully copies a bad deploy, a mass delete or a corrupting bug to the standby within seconds.
| Failure | What protects you | What does not |
|---|---|---|
| One RegionServer dies | WAL splitting and region reassignment; region replicas for read availability | Nothing at the DR layer is needed |
| Whole cluster or site lost | A standby cluster fed by replication, or restored from exported snapshots | Snapshots stored on the lost cluster's HDFS |
| HDFS or metadata corruption | Standby cluster, exported snapshots | Region replicas, which share the same HDFS |
| Logical error: bad writes, deletes, dropped table | Snapshots and archived WALs taken before the error | Replication, which has already applied the error |
A complete design therefore has two independent halves: a standby for site loss and a backup program for logical errors.
The reference topology
The standby is a separate cluster with its own HDFS, its own ZooKeeper quorum and, ideally, its own authentication and name services. Data reaches it by replication of WAL edits, and a second path exports snapshots to an object store that neither cluster depends on. Clients reach HBase through a name you control, such as a DNS record or a service-discovery entry holding the ZooKeeper quorum or connection URI, so failover is a configuration change rather than a code deployment.
Everything besides table data
Replication copies cells, not the environment those cells need. A standby missing any of the following fails on the first client request.
- Schema. Replication does not propagate DDL. Every table, column family and setting must exist on the standby, including
REPLICATION_SCOPE, TTLs, versions and compression. Check parity with a script, not by eye. - Namespaces, ACLs, quotas and RegionServer groups. Security and multi-tenancy configuration lives outside user tables and must be applied separately.
- Coprocessors and custom filters. Their jars must be deployed on the standby at the same version, or regions fail to open.
- Layered metadata. If you use Apache Phoenix, its SYSTEM tables and index tables must be part of the plan; replicating data tables without them gives you tables the query layer cannot read.
- Authentication. With Kerberos, the standby needs reachable KDCs and the service principals and keytabs your clients use. A KDC that only exists at the primary site turns a site failure into a total outage.
- Capacity. A standby sized only as a replication sink cannot serve production load; size it for failover or pre-agree which tables are served.
Bootstrapping an asynchronous standby
Replication only ships edits written after a peer exists. Existing data has to be copied separately, and the order of steps decides whether you end up with a gap. The safe order is: create the peer so the primary starts retaining WALs for it, keep the peer disabled, copy a snapshot, then enable the peer so the queued edits are applied on top.
# On the primary (HBase shell). Create the peer disabled so WALs queue for it.
# HBase 2 shells accept STATE on add_peer; on older shells run disable_peer immediately.
add_peer '1', CLUSTER_KEY => "dr-zk1,dr-zk2,dr-zk3:2181:/hbase", STATE => "DISABLED"
alter 'orders', {NAME => 'd', REPLICATION_SCOPE => 1}
snapshot 'orders', 'orders_bootstrap'
# Copy the snapshot to the DR cluster's HDFS (run from the primary).
hbase org.apache.hadoop.hbase.snapshot.ExportSnapshot \
-snapshot orders_bootstrap -copy-to hdfs://dr-nn:8020/hbase -mappers 16
# On the DR cluster (HBase shell): materialize the table from the snapshot.
clone_snapshot 'orders_bootstrap', 'orders'
# Back on the primary: start shipping the queued edits.
enable_peer '1'Edits written between the snapshot and enable_peer are shipped and applied on top of the clone. Some of them are already in the snapshot, which is harmless: replicated cells keep their original timestamps, so applying the same put or delete twice produces the same result. Check help 'add_peer' on your version before relying on the STATE option.
What RPO async replication really gives you
With asynchronous replication, a write is acknowledged as soon as it reaches the primary's WAL. It is shipped later. If the primary site disappears, everything not yet applied at the standby is lost, at least until the primary's disks can be read again. Your RPO is therefore the replication lag at the moment of failure, and lag grows exactly when you least want it to: during load spikes, network trouble and RegionServer restarts.
Watch two source-side metrics per RegionServer: the age of the last shipped operation and the size of the WAL queue. Alert on the age crossing a fraction of your RPO, and treat a queue that keeps growing as a slow incident even when the age looks fine.
# Alert when any RegionServer's replication age threatens the RPO.
import requests
RPO_MS = 30_000
for host in REGIONSERVERS:
beans = requests.get(f"http://{host}:16030/jmx", timeout=5).json()["beans"]
for b in beans:
if b.get("name", "").startswith("Hadoop:service=HBase,name=RegionServer,sub=Replication"):
for key, value in b.items():
if key.endswith("ageOfLastShippedOp") and value > RPO_MS / 3:
page(f"{host} {key}={value}ms, RPO budget {RPO_MS}ms")Metric names differ slightly between versions, so match on suffixes as above and confirm against your own /jmx output.
Synchronous replication
For tables that cannot lose acknowledged writes, HBase offers synchronous replication between two clusters. A cluster in the ACTIVE state writes every WAL entry both locally and to a remote WAL directory on the peer cluster's HDFS before acknowledging, while normal replication applies the edits to the peer's tables. The peer is in the STANDBY state, in which it rejects both reads and writes from clients. DOWNGRADE_ACTIVE is the state of a cluster that accepts writes but no longer writes a remote WAL, used when the other side is unavailable.
# Same peer id on BOTH clusters; table-level replication only.
# On cluster A:
add_peer '1', CLUSTER_KEY => "zk-b:2181:/hbase",
REMOTE_WAL_DIR => "hdfs://nn-b:8020/hbase/remoteWALs",
TABLE_CFS => {"ledger" => []}
# On cluster B:
add_peer '1', CLUSTER_KEY => "zk-a:2181:/hbase",
REMOTE_WAL_DIR => "hdfs://nn-a:8020/hbase/remoteWALs",
TABLE_CFS => {"ledger" => []}
transit_peer_sync_replication_state '1', 'STANDBY' # on B
transit_peer_sync_replication_state '1', 'ACTIVE' # on AThe costs are concrete. Every write now includes a round trip to the remote HDFS, so write latency rises with the distance between sites and the remote cluster's health becomes part of your write path. The standby cannot serve reads, so it is pure insurance, not a read replica. Only table-level configuration is supported, and the peer id must be the same on both clusters. Use it for the small set of tables whose RPO truly is zero, and keep the rest on asynchronous replication.
The documented procedures are short. If the standby fails, move the active cluster to DOWNGRADE_ACTIVE so writes continue and edits queue asynchronously until the standby returns. If the active cluster fails, move the standby to DOWNGRADE_ACTIVE, redirect clients to it, and when the old active recovers, move it to STANDBY before anything else so it cannot accept writes that diverge.
Failover runbook for an async standby, with a worked example
Suppose the primary site loses power. The orders table takes 50,000 writes per second and replication age was 4 seconds at the last scrape, so roughly 200,000 recent edits may be missing at the standby. The RTO is 30 minutes. The runbook below is executed by an incident commander, with each step recorded in the incident channel.
- Declare and decide. Confirm the primary is unreachable from several vantage points. Failing over on a network partition, while the primary still accepts writes, creates two diverging copies.
- Fence the primary. Stop writers from reaching it: point the client endpoint at nothing, or revoke write permission, so a half-alive primary cannot take writes that the standby will never see.
- Record the cut-over point. Note the last replication age and the wall-clock time. This bounds the window of possibly lost edits, here roughly the 4 seconds before the outage.
- Prepare the standby. Make sure it has no enabled peer pointing back at the dead primary, apply any schema or ACL drift, and scale up RegionServers if it was sized as a sink.
- Redirect clients. Change the DNS record or service entry and restart or refresh client connections. Measure the error rate and latency before announcing recovery.
When the old primary's disks are readable again, its WALs still contain the unshipped edits. Replay the time window into the standby with WALPlayer, limited to the tables and time range at risk. Because cells carry their original timestamps, replaying an edit that did arrive is harmless, and older replayed cells do not mask newer standby writes unless the old primary's clock ran ahead. The real risk is decisions made on the standby without the missing data, so review critical rows such as balances rather than replaying blindly.
# Replay WALs copied from the old primary for the at-risk window (milliseconds since epoch).
hbase org.apache.hadoop.hbase.mapreduce.WALPlayer \
-Dwal.start.time=1790839500000 -Dwal.end.time=1790839520000 \
hdfs://dr-nn:8020/recovered/oldwals orders
Failback without divergence
Failback is a second failover in the opposite direction and deserves the same care. The old primary is stale and may hold edits that never shipped. Treat it as a new standby: set up replication from the current active to it, bootstrap it from a fresh snapshot if it has missed too much, verify it, then schedule a planned switch during low traffic, fencing the current active before redirecting clients.
Running both clusters active at once, with replication in both directions, removes the switch but changes the consistency model. Concurrent writes to the same cell on both sides are resolved by timestamp, so the result depends on clock skew between sites. That is acceptable for idempotent, single-writer-per-key workloads and wrong for counters, balances or anything read-modify-write.
Verify continuously and drill
A standby you have never compared with the primary is a hypothesis. VerifyReplication runs a MapReduce job that compares rows between a table and its peer over a time range and reports good and bad rows as job counters. Run it on a schedule for recent time windows, excluding the last few minutes so in-flight edits do not show up as mismatches.
hbase org.apache.hadoop.hbase.mapreduce.replication.VerifyReplication \
--starttime=1790830000000 --endtime=1790833600000 1 ordersThen drill. Once a quarter, fail a non-critical table over for real, serve traffic from the standby, and fail back, timing every step. Drills typically find a missing coprocessor jar, a client caching the old quorum, or an expired keytab at the DR site.
Failure modes
- Peer disabled and forgotten. WALs pile up on the primary's HDFS while the standby silently ages. Alert on peer state as well as lag.
- Schema drift. A column family added only on the primary makes replicated edits fail at the sink, and the queue grows.
- Shared dependency. Both sites use one KDC, one DNS zone or one object-store bucket, and the disaster takes it out too.
- Split brain. Failover during a partition leaves two writable copies; reconciling them is manual and slow.
What to do next
- Write an RPO and RTO for every table with its owner, and assign each table to a tier.
- Graph replication age and WAL queue size per RegionServer at peak, and alert at a third of each table's RPO.
- Script a schema, ACL, quota and coprocessor parity check between primary and standby, and run it daily.
- List every dependency of the standby, such as KDC, DNS and object storage, and remove any that live only at the primary site.
- Schedule VerifyReplication for recent windows on critical tables.
- Run a timed failover and failback drill on a non-critical table, then update the runbook with what broke.
- Confirm that a separate backup program covers logical errors, independent of the standby.