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.

Advertisement

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

Snapshot plus replication: the bulk copy and the live tail meet at the snapshot pointSource clusterTarget clusterWALs for the table start queueingmanifest of HFiles at time Tedits after T, held for the peertable as of Tcatch up to liveVerifyReplication, SyncTableHFiles (MapReduce copy)WAL edits (replication)the queue grows until step 4:watch disk on the source
Ordering is the whole trick: the peer exists before the snapshot, so every edit after the snapshot point is queued for the target, and nothing falls in a gap between the copy and the tail.

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.

Advertisement

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 clicks

Expect 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

PhaseDurationWhat was watched
Inventory, codecs, coprocessor jars, Kerberos trust1 weektest region opens on target with a small table
Peer added and disabled; snapshotminutesoldWALs growth rate: ~3 TB/day raw
ExportSnapshot, 64 mappers, capped at 5 Gbps total~18 hourslink utilisation, source HDFS free space
Clone, enable peer, catch-up~6 hoursageOfLastShippedOp from ~24 h to under 1 s
VerifyReplication + SyncTable dry run~10 hoursBADROWS explained: TTL expiry only
Readers moved service by service3 daysp99 latency on target vs source
Write pause, final verify, flip, reverse peer12 minuteserror 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

FailureSymptomPrevention or fix
WAL backlog fills the sourceHDFS space alarms while the peer is disabledcompute backlog from write rate x export time; add space first
Missing codec or coprocessor on targetregions stuck opening, replication stalls on one serverpre-flight test table with every family option
Bulk loads not replicatedrows missing only for some ingestion jobsenable bulk-load replication or re-run loads on the target
Snapshot taken before the peerunexplained gaps in verificationalways create and disable the peer first
Uncapped exportapplication latency spikes, replication lag growsset -mappers and -bandwidth together
No reverse pathrollback loses writes made on the targetadd the reverse peer before unpausing writers

What to do next

  1. Inventory every table's families, codecs, coprocessors, ACLs and ingest paths, and open a small test table with the same options on the target.
  2. Compute export time from table size and usable bandwidth, and the source disk needed to hold write-ahead logs for that long.
  3. Rehearse the full sequence on one small table: disabled peer, snapshot, ExportSnapshot, clone, enable, verify.
  4. Run VerifyReplication and a SyncTable dry run, and write down an explanation for every mismatch.
  5. Move clients to a configuration-driven cluster address so cutover and rollback are both a config change.
  6. Plan the cutover with a reverse peer in place and a date for removing the source.
Key takeaway: A safe HBase migration is ordering plus evidence: create a disabled replication peer before the snapshot so no write falls in a gap, move the bulk with a capped ExportSnapshot, let replication replay the tail, prove equality with VerifyReplication and a SyncTable dry run, and cut over with a reverse peer already in place so rollback never loses data.