Migrating an HBase table means moving tens of terabytes that never stop changing, to a cluster that may run a different version, in a different network, sometimes on a different platform, while applications keep reading and writing. The copy is the easy part. The hard parts are making sure no write made during the copy is lost, proving that the two sides really match, and having a way back if the new cluster misbehaves after cutover.
This article is about moving data between clusters: a data-centre exit, a move from on-premises to a cloud-hosted cluster, a change of distribution, or a move to Bigtable. Upgrading a cluster's version where it stands is a different job, covered in HBase upgrade in depth. We will build the standard pattern from first principles (a replication peer, a snapshot, an export, a catch-up and a verified cutover), work through a 40 TB example, and list the failures that turn a weekend migration into a month.
What a correct migration guarantees
Define success before choosing tools. A correct migration guarantees four things. No lost writes: every mutation acknowledged by the source before cutover is present on the target. Verified equality: you have evidence, not hope, that the tables match at a known point. Bounded downtime: writes pause for minutes, not hours, or not at all. Reversibility: for some period after cutover you can return to the source without losing writes made on the target.
Two HBase properties make this achievable. Every cell carries a timestamp, and writing the same cell (same row, column and timestamp) with the same value twice is harmless; that makes replaying edits idempotent. And HBase can both capture a table as a set of immutable files (a snapshot) and stream every edit from the write-ahead log to another cluster (replication). The pattern combines them: the snapshot moves the bulk, replication moves everything that happened after it.
The pattern: a disabled peer, then a snapshot
The order of operations is what guarantees no lost writes. First, add a replication peer pointing at the target and leave it disabled. A disabled peer still has its write-ahead logs tracked: the source keeps every log file that contains edits for replicated column families instead of deleting it. Second, take a snapshot. Every edit is now either in the snapshot, in the queued logs, or in both. Third, export the snapshot and create the table from it on the target. Fourth, enable the peer; the source ships the backlog, and the target applies it on top of the snapshot. Edits that were already in the snapshot are simply written again with the same timestamps, which changes nothing.
If you snapshot first and add the peer afterwards, edits made in between are in neither place, and you will not find out until verification, if at all. Replication itself (sources, sinks, queues and ordering) is described in HBase replication architecture, and how a snapshot captures a table without copying it in HBase snapshots.
Pre-flight inventory
Most failed migrations fail on something the target lacks, not on the copy. Before touching data, inventory the source and check each item on the target:
- Tables and schemas: namespaces, column families, versions, TTLs,
KEEP_DELETED_CELLS, block encoding, and region split points. - Compression codecs: a family compressed with Snappy or ZSTD needs working native libraries on every target RegionServer, or regions fail to open.
- Code on the classpath: coprocessors, custom filters and custom comparators must be deployed on the target before any region that uses them opens.
- Security: ACLs or Ranger policies, cell labels, Kerberos principals and, across realms, a cross-realm trust so the export job and replication can authenticate.
- Ingest paths: bulk-loaded HFiles bypass the write-ahead log, so replication does not see them unless bulk-load replication is enabled on both sides (see bulk load architecture).
- Client-set timestamps: applications that set their own timestamps can overwrite newer data with older data during replay; find them now.
Then do the bandwidth arithmetic, because it sets the schedule. Copying 60 TB over a 10 Gbps link that you may use at 60 percent is 60 x 10^12 x 8 bits divided by 6 x 10^9 bits per second, which is 80,000 seconds or about 22 hours. The peer is disabled for that whole time, so the source must hold 22 hours of write-ahead logs plus a margin. At 5 TB of writes a day, that is several terabytes of extra HDFS space, multiplied by the replication factor.
Step 1 and 2: create the peer, take the snapshot
# On the SOURCE cluster, in the hbase shell.
# Mark the families to replicate (REPLICATION_SCOPE 1 = replicate).
alter 'clicks', {NAME => 'e', REPLICATION_SCOPE => 1}
# Add a peer for the target; restrict it to this table, then disable it at once.
add_peer '7', CLUSTER_KEY => "tzk1,tzk2,tzk3:2181:/hbase", TABLE_CFS => { "clicks" => ["e"] }
disable_peer '7'
list_peers
# Take the snapshot only after the peer is in place.
snapshot 'clicks', 'clicks-mig-20261001'Adding and immediately disabling the peer leaves a short window where the source may try to ship edits. That is safe: the target table does not exist yet, so shipping fails and is retried later; nothing is dropped. Some shell versions accept a state option on add_peer to create the peer disabled in one step; check help 'add_peer' on your version before relying on it.
Step 3: export and clone
ExportSnapshot is a MapReduce job that copies the snapshot's HFiles and manifest to another filesystem. It does not go through the RegionServers, so it does not load the source's read path, only HDFS and the network. Run it with an explicit number of mappers and a bandwidth cap, because an uncapped export can saturate the link that replication and your applications also use.
# Run from a node with access to both filesystems.
hbase org.apache.hadoop.hbase.snapshot.ExportSnapshot \
-snapshot clicks-mig-20261001 \
-copy-to hdfs://target-nn:8020/hbase \
-mappers 64 \
-bandwidth 10 # MB/s per mapper: 64 x 10 = ~5 Gbps
# Then, on the TARGET cluster:
# hbase shell> clone_snapshot 'clicks-mig-20261001', 'clicks'The bandwidth option limits each mapper, so the total is roughly mappers times the cap; size both together against your link. Clone the snapshot on the target to create the table with the source's split points, which avoids a storm of splits. After cloning, the target table's files are links to the snapshot's files until compactions rewrite them. The cleaner tracks those references, so deleting the snapshot does not break the table, but keep it until cutover as a known-good restore point.
Step 4: catch up
Enable the peer and watch the backlog drain. The source ships queued edits in batches to RegionServers on the target, which apply them as ordinary writes. Two metrics matter: the age of the last shipped operation, which tells you how far behind the target is in time, and the size of the log queue, which tells you how many files are still waiting. Both should fall steadily. status 'replication' in the shell shows them per RegionServer.
enable_peer '7'
status 'replication'
# Watch per RegionServer: ageOfLastShippedOp falling toward 0, sizeOfLogQueue shrinking.Catch-up speed depends on how many RegionServers ship in parallel and how fast the target can absorb writes. A target that has just been cloned has cold caches and will compact heavily, so expect the first hours to be slower. If the age stops falling, look for a single RegionServer with a growing queue; that usually means one region on the target is failing writes, for example because a coprocessor is missing.
Step 5: verify
Verification needs two tools because they answer different questions. VerifyReplication is a MapReduce job that scans the source table and compares each row with the peer's copy, reporting matching and mismatching rows in its counters (GOODROWS and BADROWS). Give it a time range that ends a few minutes in the past, so rows still in flight are not counted as mismatches.
HashTable and SyncTable compare the whole table efficiently. HashTable scans the source and writes a hash for each batch of rows; SyncTable scans the target, hashes the same batches, and only looks at individual cells where hashes differ. Run SyncTable with --dryrun=true first, so it reports differences instead of fixing them; a difference you do not understand is a bug to investigate, not something to overwrite.
# Row-level comparison against peer 7 for a closed time window (ms since epoch).
hbase org.apache.hadoop.hbase.mapreduce.replication.VerifyReplication \
--starttime=1790726400000 --endtime=1790812800000 7 clicks
# Whole-table hash comparison: hash on the source, diff on the target.
hbase org.apache.hadoop.hbase.mapreduce.HashTable \
--batchsize=32000 --numhashfiles=64 clicks /migration/hashes/clicks
hbase org.apache.hadoop.hbase.mapreduce.SyncTable --dryrun=true \
--sourcezkcluster=szk1,szk2,szk3:2181:/hbase \
hdfs://source-nn:8020/migration/hashes/clicks clicks clicksExpect some benign differences and explain each one. TTL expiry happens at different moments on the two clusters; version limits can keep different old versions until compaction; delete markers behave differently when KEEP_DELETED_CELLS differs. Add an application-level check too: read a sample of real keys through your own client code on both sides and compare the decoded results.
Step 6: a reversible cutover
The safest cutover pauses writes briefly. Stop or pause writers (or put the application in read-only mode), wait until replication lag reaches zero, run a final VerifyReplication over the last few minutes, then switch clients to the target by changing the ZooKeeper quorum or connection configuration they load. Before unpausing writers, add a reverse peer from the target to the source. From that moment, writes on the new cluster also flow back to the old one, and rollback is another short pause and a configuration switch in the other direction.
Readers can move before writers, one service at a time, because the target is kept current by replication; that spreads the risk. Dual writing from the application, where every write goes to both clusters, is the alternative when you cannot pause at all, but it moves consistency into your code: a failure on one side needs a retry queue and a reconciliation job. Prefer replication when you can.
// Clients read their cluster from a config service so cutover is a config change.
Configuration conf = HBaseConfiguration.create();
conf.set("hbase.zookeeper.quorum", config.get("clicks.cluster.quorum")); // flip here
try (Connection conn = ConnectionFactory.createConnection(conf);
Table t = conn.getTable(TableName.valueOf("clicks"))) {
// unchanged application code
}
Worked example: 40 TB of clickstream
| Phase | Duration | What was watched |
|---|---|---|
| Inventory, codecs, coprocessor jars, Kerberos trust | 1 week | test region opens on target with a small table |
| Peer added and disabled; snapshot | minutes | oldWALs growth rate: ~3 TB/day raw |
| ExportSnapshot, 64 mappers, capped at 5 Gbps total | ~18 hours | link utilisation, source HDFS free space |
| Clone, enable peer, catch-up | ~6 hours | ageOfLastShippedOp from ~24 h to under 1 s |
| VerifyReplication + SyncTable dry run | ~10 hours | BADROWS explained: TTL expiry only |
| Readers moved service by service | 3 days | p99 latency on target vs source |
| Write pause, final verify, flip, reverse peer | 12 minutes | error rate, replication lag both ways |
The team kept the reverse peer and the source cluster for two weeks after cutover, then removed both peers, dropped the snapshot and decommissioned the source. The only surprise was a nightly bulk-load job that had been writing HFiles directly; it had to be pointed at the new cluster before the flip, because replication had never carried its data.
Migrating to Bigtable
Bigtable speaks the HBase client API, so many applications move with a dependency and configuration change, but the server side is different. Google provides tools for each step of the same pattern: a Schema Translation Tool that reads HBase table definitions (families, garbage-collection policies, split points) and creates matching Bigtable tables; snapshot export to Cloud Storage followed by an import into Bigtable using Dataflow; an HBase Bigtable Replication Library that plugs into HBase replication so live writes follow the bulk import; and a Migration Validation Tool for comparing the two. Coprocessors and custom server-side filters do not carry over, so logic that ran inside RegionServers must move into the application or a pipeline. See Bigtable architecture for how the storage model differs.
Failure modes
| Failure | Symptom | Prevention or fix |
|---|---|---|
| WAL backlog fills the source | HDFS space alarms while the peer is disabled | compute backlog from write rate x export time; add space first |
| Missing codec or coprocessor on target | regions stuck opening, replication stalls on one server | pre-flight test table with every family option |
| Bulk loads not replicated | rows missing only for some ingestion jobs | enable bulk-load replication or re-run loads on the target |
| Snapshot taken before the peer | unexplained gaps in verification | always create and disable the peer first |
| Uncapped export | application latency spikes, replication lag grows | set -mappers and -bandwidth together |
| No reverse path | rollback loses writes made on the target | add the reverse peer before unpausing writers |
What to do next
- Inventory every table's families, codecs, coprocessors, ACLs and ingest paths, and open a small test table with the same options on the target.
- Compute export time from table size and usable bandwidth, and the source disk needed to hold write-ahead logs for that long.
- Rehearse the full sequence on one small table: disabled peer, snapshot, ExportSnapshot, clone, enable, verify.
- Run VerifyReplication and a SyncTable dry run, and write down an explanation for every mismatch.
- Move clients to a configuration-driven cluster address so cutover and rollback are both a config change.
- Plan the cutover with a reverse peer in place and a date for removing the source.