HBase and Apache Iceberg solve opposite problems. HBase serves millisecond reads and writes on individual rows keyed by a byte array, and it keeps the latest state of each row mutable in place. Iceberg turns a directory of Parquet or ORC files on object storage into a table with atomic commits, snapshots, schema evolution and hidden partitioning, so engines such as Spark, Trino and Flink can scan billions of rows consistently. Many platforms need both: HBase holds the operational state an application reads and writes, and Iceberg holds the analytical copy that reporting, feature engineering and model training scan.
The first fact to get straight is that there is no native integration. Iceberg has no HBase-backed catalog and cannot store its data or metadata in HBase, and HBase has no Iceberg table format. Iceberg's catalogs are services such as the Hive Metastore, a REST catalog, JDBC, Nessie or AWS Glue. Integration is data movement, and it is usually written in Spark or Flink. This article covers three movement patterns, the mapping they share, a worked example and the failure modes that make the copies quietly disagree.
Why pair a serving store with a lakehouse table
Consider a customer-profile service. The application looks up a profile by customer id thousands of times a second and updates it whenever the customer changes tier or email. That workload suits HBase: rows are addressed by key, writes go to a memstore and a write-ahead log, and a read costs a few block lookups. Now analysts want to know how many customers moved from silver to gold last quarter, and a training job wants every profile joined to a year of orders. Scanning HBase for each question competes with production traffic.
Iceberg is the right home for those questions. Data sits in columnar files, scans read only the needed columns, and each commit produces a snapshot you can query later, which is what makes a training set reproducible. The cost is latency: Iceberg commits are batch operations, typically every few minutes at best, and each commit adds files that later need compaction. So the architecture is a division of labour. HBase owns the current truth for serving. Iceberg owns history and analysis. The integration decides how fresh the Iceberg copy is and how expensive it is to keep honest.
Mapping cells to columns
Every pattern needs the same translation, and most correctness bugs live here rather than in the pipeline plumbing. An HBase cell is a row key, family, qualifier, timestamp, type and value, all as raw bytes. An Iceberg row has typed columns. Decide the following explicitly and write it down:
- Row key. If the key is salted or composite, for example a one-byte salt followed by the customer id, decode it into real columns and drop the salt. Keep the raw key as a binary column too if you will ever write back to HBase.
- Value encoding. HBase stores whatever bytes the writer chose. A long written with
Bytes.toBytes(long)is eight big-endian bytes, while a value written as a string is UTF-8 text. Only the application knows which, so the mapping must be a table of qualifier, encoding and Iceberg type. - Dynamic qualifiers. Families where qualifiers are data, such as one column per day, do not fit fixed columns. Map them to an Iceberg
map<string, string>column, or pivot them into a child table. - Versions. HBase may keep several timestamped versions per cell. Most pipelines take the newest and also store its timestamp as
hbase_ts, which later serves as the ordering key for change data. - Deletes. HBase deletes are markers, not removals. A client row delete is written as family delete markers for every family, a column delete as a column marker. Map a family delete covering every replicated family to an Iceberg row delete, and column deletes to an update that sets that column to null.
- TTL. Cells that expire through a family TTL vanish during compaction without any delete event. If the lake must mirror that, apply the same retention as a scheduled Iceberg delete.
Pattern A: snapshot export
The simplest pattern, and the right one for the initial load, backfills and periodic full reconciliation, reads an HBase snapshot directly. A snapshot is a cheap metadata operation that references the table's current HFiles. TableSnapshotInputFormat then lets Spark read those files straight from HDFS or object storage, bypassing the RegionServers entirely, so the export does not compete with serving traffic. See HBase Snapshots, in depth for how snapshots and their cleanup work.
# HBase shell: take the snapshot (flushes memstores first by default)
snapshot 'profile', 'profile_20261004'// Spark (Scala): read the snapshot's HFiles and write an Iceberg table
val conf = HBaseConfiguration.create()
val job = Job.getInstance(conf)
// restoreDir: same filesystem as hbase.rootdir, but NOT under it
TableSnapshotInputFormat.setInput(job, "profile_20261004", new Path("/tmp/snap_restore/profile"))
val rows = spark.sparkContext.newAPIHadoopRDD(job.getConfiguration,
classOf[TableSnapshotInputFormat], classOf[ImmutableBytesWritable], classOf[Result])
.map { case (_, r) =>
val p = Bytes.toBytes("p")
def cell(q: String) = Option(r.getColumnLatestCell(p, Bytes.toBytes(q)))
val key = r.getRow
(Bytes.toString(key, 1, key.length - 1), // drop 1-byte salt
cell("email").map(c => Bytes.toString(CellUtil.cloneValue(c))).orNull,
cell("tier").map(c => Bytes.toString(CellUtil.cloneValue(c))).orNull,
cell("ltv").map(c => java.lang.Double.valueOf(Bytes.toDouble(CellUtil.cloneValue(c)))).orNull,
new java.sql.Timestamp(r.rawCells().map(_.getTimestamp).max))
}
rows.toDF("customer_id", "email", "tier", "ltv", "hbase_ts")
.withColumn("deleted", lit(false)) // tombstone flag, see pattern B
.writeTo("lake.crm.customer_profile")
.overwritePartitions()Three operational notes. An online snapshot is consistent per row, but it is not a single global instant across all regions, so treat it as row-consistent rather than as a transaction. The restore directory must not sit under the HBase root directory, and you should delete the snapshot and the restore directory when the job finishes, or the archived HFiles it references are never cleaned up. And because the job reads files, not servers, its parallelism follows regions: one input split per region by default, so a table with twelve huge regions gives you twelve slow tasks.
Pattern B: change data capture into MERGE
To keep Iceberg minutes rather than a day behind, stream changes. HBase replication ships write-ahead-log edits for column families with REPLICATION_SCOPE => 1 to peers, and a peer does not have to be another HBase cluster: add_peer accepts a custom endpoint class, and the hbase-connectors project includes a Kafka proxy built on this mechanism. Evaluate its maturity for your version, or write a small endpoint that publishes each edit as a record keyed by row key. The mechanics, including serial replication for per-key ordering, are in HBase Replication, in depth.
On the lake side, a Spark Structured Streaming job reads the topic and applies each micro-batch with MERGE INTO. Two rules are not optional. First, collapse each micro-batch to one row per key before merging, because Spark refuses a MERGE in which several source rows match the same target row. Second, compare timestamps so that a replayed or late event cannot overwrite newer state, because replication is at-least-once. That second rule has a trap: a real DELETE removes the row and its timestamp, so a late, older put would fall through to the insert clause and resurrect it. Keep deletes as tombstones instead:
from pyspark.sql import Window, functions as F
def apply_batch(batch, batch_id):
latest = (batch
.withColumn("rn", F.row_number().over(
Window.partitionBy("customer_id").orderBy(F.col("hbase_ts").desc())))
.filter("rn = 1").drop("rn"))
latest.createOrReplaceTempView("changes")
batch.sparkSession.sql("""
MERGE INTO lake.crm.customer_profile t
USING changes s
ON t.customer_id = s.customer_id
WHEN MATCHED AND s.hbase_ts >= t.hbase_ts THEN UPDATE SET
email = s.email, tier = s.tier, ltv = s.ltv, hbase_ts = s.hbase_ts,
deleted = (s.op = 'ROW_DELETE')
WHEN NOT MATCHED THEN INSERT
(customer_id, email, tier, ltv, hbase_ts, deleted)
VALUES (s.customer_id, s.email, s.tier, s.ltv, s.hbase_ts, s.op = 'ROW_DELETE')
""")
(changes_df.writeStream.foreachBatch(apply_batch)
.option("checkpointLocation", "s3://lake/chk/customer_profile")
.trigger(processingTime="5 minutes").start())Readers query a view with WHERE NOT deleted, and a scheduled DELETE ... WHERE deleted purges tombstones older than your replay horizon, such as the Kafka retention period.
This sketch assumes each change record carries the full current row. WAL edits carry only the cells that changed, so either the endpoint reads the full row before publishing, or the merge updates only non-null columns and tracks explicit column deletes separately. Choose deliberately; the second is cheaper for HBase but makes the merge logic harder to test.
For a frequently updated table, set write.merge.mode, write.update.mode and write.delete.mode to merge-on-read on a format version 2 or later table, so each commit writes small delete files instead of rewriting whole data files. That moves cost to readers and to maintenance, which is the next section.
Keeping the Iceberg table healthy
A five-minute trigger creates nearly three hundred commits a day, each with new data and delete files. Without maintenance, scans slow down and metadata grows without limit. Schedule these Iceberg procedures:
CALL lake.system.rewrite_data_files(table => 'crm.customer_profile');
CALL lake.system.rewrite_position_delete_files(table => 'crm.customer_profile');
CALL lake.system.expire_snapshots(table => 'crm.customer_profile',
older_than => TIMESTAMP '2026-09-27 00:00:00');
CALL lake.system.remove_orphan_files(table => 'crm.customer_profile');Run compaction separately from the streaming merge and expect occasional commit conflicts, which Iceberg retries. Keep snapshots long enough for any training run that pins one, because an expired snapshot cannot be time-travelled to. Spark and Iceberg covers the procedures and commit model in more detail.
Pattern C: lake results back into HBase
The third pattern runs the other way. Features, scores or aggregates computed in the lake often need millisecond lookup by key, which is HBase's job. Writing millions of rows through the client API is slow and loads the RegionServers; generating HFiles and bulk loading them is far cheaper, as HBase Bulk Load, in depth explains. Read the Iceberg table at a pinned snapshot so a retried job loads exactly the same data:
snap = spark.sql("SELECT snapshot_id, committed_at FROM lake.crm.customer_scores.snapshots "
"ORDER BY committed_at DESC LIMIT 1").first()
scores = spark.read.option("snapshot-id", snap.snapshot_id).table("lake.crm.customer_scores")
# Write HFiles sorted by row key, using snap.committed_at as every cell's timestamp,
# then hand the directory to the bulk-load tool for your HBase version.Use the snapshot's commit time as the cell timestamp. Reloading the same snapshot then writes identical versions, which makes retries idempotent, and anyone can trace an HBase value back to the lake snapshot that produced it. The bulk-load tool's class name has changed across HBase 2.x releases, so use the one documented for your version. The hbase-spark module also offers a bulk-load helper, described in HBase and Spark Connector, in depth.
Worked example: the profile service end to end
Putting it together for the profile service: on day one, snapshot profile and run pattern A into an Iceberg table partitioned by bucket(16, customer_id), with hbase_ts on every row. Record the snapshot's creation time. Start the replication peer, then start the streaming merge from a Kafka offset before that time. Replaying the overlap is harmless, because the timestamp guard and tombstones reject older events. Nightly, compare the two sides per bucket on row count and a hash of the key and hbase_ts, using a fresh snapshot export on the HBase side:
SELECT b, count(*) AS n, bit_xor(xxhash64(customer_id, hbase_ts)) AS h
FROM (SELECT pmod(xxhash64(customer_id), 16) AS b, customer_id, hbase_ts
FROM lake.crm.customer_profile WHERE NOT deleted)
GROUP BY bRun the same query over the export and diff the sixteen rows. A mismatched bucket narrows the search to one sixteenth of the keys; repair it by re-merging that bucket's rows from the export.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| MERGE fails: one target row matched several source rows | Micro-batch holds several changes per key | Deduplicate to the newest change per key before merging |
| Old values reappear after a restart | Replayed events overwrite newer rows | Guard with a timestamp comparison and keep tombstones until the replay horizon passes |
| Deleted HBase rows persist in the lake | Delete markers not mapped, or TTL expiry | Map family-wide delete markers to DELETE; mirror TTL as a scheduled delete |
| Garbage numbers in Iceberg | Bytes decoded with the wrong encoding | Maintain an explicit qualifier-to-encoding mapping and test it |
| Archive directory keeps growing | Snapshots or restore dirs never deleted | Delete both at the end of each export job |
| Lake scans slow down week by week | Small files and delete files accumulate | Schedule compaction, delete-file rewrite and snapshot expiry |
| HBase rows never reach the lake | SKIP_WAL writes and bulk loads bypass WAL replication by default | Avoid SKIP_WAL on replicated tables; re-export bulk-loaded data; rely on reconciliation |
| Replication lag climbs | Endpoint or Kafka slower than write rate | Monitor the peer's queue size, and scale the endpoint and topic partitions |
Trade-offs and alternatives
Pattern A alone, run nightly, is the cheapest and simplest, and is often enough when analysts accept a day of staleness. It also gives you a repair tool for any other pattern. Pattern B buys freshness at the price of a replication peer, a Kafka topic, a streaming job and compaction, so adopt it only for tables where minutes matter. Pattern C is worth it whenever lake results need key lookups, which Iceberg scans cannot serve.
Two alternatives deserve a look. If you already run Apache Phoenix, its SQL view of HBase can simplify type decoding. And if the main need is SQL over HBase rather than a separate copy, the Hive storage handler in HBase and Hive Integration, in depth queries HBase in place, at the cost of load on the RegionServers.
What to do next
- Write the mapping document first: row key decoding, every qualifier's encoding and Iceberg type, version and delete rules, and TTL handling.
- Run pattern A against a staging snapshot and compare counts and sampled rows with HBase before building anything streaming.
- Add
hbase_tsand adeletedflag to the Iceberg schema now; every later pattern depends on them for ordering and deletes. - If minutes of freshness matter, enable replication scope on the needed families, add a peer and build the deduplicating, timestamp-guarded merge.
- Schedule compaction, delete-file rewrites, snapshot expiry and orphan-file removal before the table goes live.
- Add the nightly per-bucket reconciliation and alert on any mismatched bucket.
- For lake-to-HBase loads, pin the Iceberg snapshot id and use its commit time as the cell timestamp.