Apache HBase stores sparse, wide tables sorted by row key and serves single-row reads and writes in milliseconds. Apache Spark scans and transforms large datasets in parallel. Teams that run both end up wanting to join them: enrich a Spark job with lookups against HBase, run SQL over an HBase table, or write the output of a batch job into HBase without hammering the RegionServers. The HBase Spark connector is the Apache-maintained bridge for that.
This article explains how the connector maps an HBase table to a Spark DataFrame, how a SQL predicate becomes a scan range or a server-side filter, what the lower-level HBaseContext API gives you, and when to skip RPC writes entirely and bulk load HFiles. It finishes with a worked example, the failure modes that show up in production, and a checklist. Option names and method signatures below were checked against the connector source.
How Spark and HBase connect
What the connector is
The connector lives in the separate apache/hbase-connectors repository, not in HBase itself, and ships as the hbase-spark module. Its 1.0.1 release made Spark 3 (built against 3.1.2), Scala 2.12 and Hadoop 3 the default build, with HBase 2.4 as the target. Because Spark, Scala and HBase versions move independently, most teams build it for their exact stack using the Maven properties the README documents:
mvn -Dspark.version=3.1.2 -Dscala.version=2.12.10 -Dscala.binary.version=2.12 \
-Dhadoop-three.version=3.2.0 -Dhbase.version=2.4.8 clean installIt offers two layers. The first is a Spark SQL data source, registered under the class name org.apache.hadoop.hbase.spark, which lets you read and write DataFrames. The second is HBaseContext, a Scala/Java API that hands each partition an HBase Connection and wraps common patterns (bulk put, bulk get, bulk delete, scan to RDD, bulk load). The data source is built on top of HBaseContext.
Two setup steps are easy to miss. The Spark side needs hbase-site.xml in $SPARK_CONF_DIR so executors can find ZooKeeper and the cluster. And if you want column predicates evaluated on the RegionServers (the default), the RegionServers need scala-library, hbase-spark and hbase-spark-protocol-shaded on their classpath, because the pushed-down filter is a connector class that runs inside the server. If you cannot touch the servers, set hbase.spark.pushdown.columnfilter to false and Spark evaluates those predicates itself.
Mapping an HBase table to a DataFrame
HBase has no schema beyond column families; every cell is bytes. The connector needs a catalog that says which bytes are which Spark column and type. There are two formats. The JSON catalog names the table, declares the row key, and maps each column to a family and qualifier; row-key fields use the pseudo-family rowkey:
val catalog = s"""{
"table": {"namespace": "default", "name": "events"},
"rowkey": "key",
"columns": {
"event_key": {"cf": "rowkey", "col": "key", "type": "string"},
"user_id": {"cf": "d", "col": "uid", "type": "string"},
"kind": {"cf": "d", "col": "kind", "type": "string"},
"amount": {"cf": "d", "col": "amt", "type": "double"},
"ts": {"cf": "d", "col": "ts", "type": "long"}
}
}"""The older, terser form is a mapping string passed as hbase.columns.mapping alongside hbase.table, for example event_key STRING :key, user_id STRING d:uid, amount DOUBLE d:amt. The :key entry marks the row key. Supported types include string, int, long, short, byte, float, double, boolean, binary and timestamp.
The default encoder writes values with HBase's Bytes.toBytes. That matters for ordering: HBase sorts row keys as unsigned bytes, and a big-endian two's-complement integer puts negative numbers after positive ones. A range predicate on a signed numeric row key can therefore select the wrong rows or force a full scan. Design row keys that sort correctly as bytes, such as zero-padded strings or offset-encoded numbers, before you rely on range pushdown.
Reading: how predicates become scans
Reading is ordinary DataFrame code. The options below are the connector's own keys:
import org.apache.hadoop.hbase.spark.datasources.HBaseTableCatalog
val df = spark.read
.format("org.apache.hadoop.hbase.spark")
.option(HBaseTableCatalog.tableCatalog, catalog)
.option("hbase.spark.use.hbasecontext", "false") // build one from hbase-site.xml
.option("hbase.spark.query.cachedrows", "1000") // Scan caching: rows per RPC
.load()
val day = df.filter($"event_key" >= "20261001#" && $"event_key" < "20261002#" && $"kind" === "refund")
day.groupBy($"user_id").sum("amount").show()Here is what happens. The driver reads the region boundaries for the table and creates one Spark partition per region. Predicates on row-key columns become scan start and stop rows, so regions outside the range are pruned and never read. An equality predicate on the full row key becomes a Get. Predicates on other columns are pushed to the RegionServers as a filter, so non-matching rows never cross the network. The connector understands EqualTo, LessThan, GreaterThan, LessThanOrEqual, GreaterThanOrEqual, StringStartsWith, IsNull, IsNotNull and their And/Or combinations. Anything else, such as IN, suffix LIKE or UDFs, is evaluated by Spark after the rows arrive.
Other read options control the scan itself: hbase.spark.query.batchsize (cells per Result), hbase.spark.query.cacheblocks (false by default, which keeps a one-off analytical scan from evicting hot data from the block cache), hbase.spark.query.timestamp, hbase.spark.query.timerange.start/.end and hbase.spark.query.maxVersions. Check the physical plan with explain() to see which filters reached the source.
Writing a DataFrame
Writing a DataFrame uses the same catalog. The connector writes through HBase's TableOutputFormat, which means ordinary Puts sent over RPC. Create the table yourself first, pre-split on boundaries taken from your real key distribution, then write:
// hbase shell, once: split points chosen from the actual date#user key space
// create 'events', 'd', SPLITS => ['20260401#', '20260701#', '20261001#']
df.write
.format("org.apache.hadoop.hbase.spark")
.option(HBaseTableCatalog.tableCatalog, catalog)
.save()The connector also has a newTable option (HBaseTableCatalog.newTable) that creates a missing table, but only when the region count you pass is greater than 3, and it computes split points evenly between regionStart and regionEnd, which default to aaaaaaa and zzzzzzz. Keys like 20261001#u001 sort below the first split, so every row lands in one region. Keep newTable for quick tests.
Pre-splitting matters. A table with one hot region sends every Put from every executor to one RegionServer until it splits, which is the classic cause of a write job that is fast for a minute and then crawls. If your keys are monotonically increasing, salt or hash a prefix so writes spread across regions.
HBaseContext: the lower-level API
When you need more control than a DataFrame gives, use HBaseContext directly. It broadcasts the HBase configuration to executors, caches one Connection per executor, and closes it for you. The core methods, with their real signatures, are bulkPut(rdd, tableName, f: T => Put), bulkGet(tableName, batchSize, rdd, makeGet, convertResult), bulkDelete(rdd, tableName, f, batchSize), hbaseRDD(tableName, scan), and foreachPartition/mapPartitions which give your function the partition iterator plus a Connection. Importing HBaseRDDFunctions._ adds the same operations as methods on any RDD:
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName}
import org.apache.hadoop.hbase.client.{Get, Put, Result}
import org.apache.hadoop.hbase.spark.HBaseContext
import org.apache.hadoop.hbase.spark.HBaseRDDFunctions._
import org.apache.hadoop.hbase.util.Bytes
val hc = new HBaseContext(spark.sparkContext, HBaseConfiguration.create())
val profiles = TableName.valueOf("user_profile")
// Enrich: look up each user id in HBase, 1,000 Gets per multi-get RPC.
val ids = spark.sparkContext.parallelize(Seq("u001", "u002", "u003"))
val tiers = ids.hbaseBulkGet[(String, String)](hc, profiles, 1000,
id => new Get(Bytes.toBytes(id)).addColumn(Bytes.toBytes("p"), Bytes.toBytes("tier")),
(r: Result) => (Bytes.toString(r.getRow),
Option(r.getValue(Bytes.toBytes("p"), Bytes.toBytes("tier"))).map(Bytes.toString).orNull))
// Write back: one Put per record, batched by the client's BufferedMutator.
tiers.hbaseBulkPut(hc, TableName.valueOf("user_tier_snapshot"), { case (id, tier) =>
new Put(Bytes.toBytes(id)).addColumn(Bytes.toBytes("s"), Bytes.toBytes("tier"),
Bytes.toBytes(Option(tier).getOrElse("none")))
})The batch size on bulkGet is the main lever: larger batches mean fewer round trips but bigger requests that can time out on a busy server. For the DataFrame path the equivalent knob is hbase.spark.bulkget.size, which defaults to 1000.
Bulk loading HFiles from Spark
RPC writes go through the full write path: WAL append, MemStore, flushes and compactions. For a one-time backfill or a nightly rebuild of hundreds of millions of rows, that is wasted work. Bulk load skips it: Spark writes HFiles, already sorted and split by region boundaries, into a staging directory, and HBase then moves them into the regions as a file operation. The general mechanics are covered in HBase bulk load; this is the Spark side:
import org.apache.hadoop.fs.Path
import org.apache.hadoop.hbase.spark.KeyFamilyQualifier
import org.apache.hadoop.hbase.tool.BulkLoadHFiles
val cf = Bytes.toBytes("d")
val staging = "hdfs:///tmp/bulk/events_20261003"
rowsRdd.hbaseBulkLoad(hc, TableName.valueOf("events"), (r: EventRow) => {
val key = Bytes.toBytes(r.eventKey)
Iterator(
(new KeyFamilyQualifier(key, cf, Bytes.toBytes("uid")), Bytes.toBytes(r.userId)),
(new KeyFamilyQualifier(key, cf, Bytes.toBytes("amt")), Bytes.toBytes(r.amount)))
}, staging)
// Hand the HFiles to the RegionServers (same as `hbase completebulkload`).
BulkLoadHFiles.create(hc.config).bulkLoad(TableName.valueOf("events"), new Path(staging))Each record expands to one cell per column, and the connector shuffles and sorts all of them by row key, family and qualifier. For tables with few columns per row, hbaseBulkLoadThinRows takes one (ByteArrayWrapper, FamiliesQualifiersValues) per row instead, which shuffles rows rather than cells and is usually cheaper. Bulk-loaded data bypasses the WAL, so it does not reach replication peers that tail the WAL unless bulk-load replication is enabled on the cluster.
Worked example: one day of events
Take an events table of 1.2 billion rows, about 120 GB, in 400 regions of roughly 300 MB each. Row keys are yyyyMMdd#userId#seq, so one day of data is a contiguous range. An analyst wants refund totals per user for 1 October.
With the catalog above, the predicate event_key >= '20261001#' AND event_key < '20261002#' becomes scan bounds. If one day is about 1.3 GB, that range touches about five regions, so Spark launches about five tasks, not 400, and reads about 1.3 GB, not 120 GB. The kind = 'refund' predicate is pushed to the RegionServers, so if refunds are 2% of events, about 26 MB crosses the network. Only the aggregation runs in Spark.
Now drop the date and filter only on user_id = 'u001'. That column is not in the row key, so there are no scan bounds: all 400 regions are scanned, and the filter runs server-side on 120 GB. The query still works, but it costs a full table scan and competes with the online workload. The fix is a data-model change, such as a second table keyed by user, not a connector setting. That is the main lesson: the connector is only as selective as your row key.
Failure modes
- ClassNotFoundException or DoNotRetryIOException on scans. Column-filter pushdown is on but the RegionServers lack the connector jars. Deploy the jars or set
hbase.spark.pushdown.columnfilter=false. - A write job that slows to a crawl. All Puts hit one region because the table was not pre-split or keys are sequential. Pre-split, salt the key, or bulk load.
- RegionServer overload during a Spark read. Hundreds of concurrent scans with large caching values compete with online traffic. Cap executor count for HBase jobs, lower
cachedrows, run heavy scans off-peak, or scan a snapshot instead of the live table. - Garbage values or wrong row counts. Catalog types do not match how the producer encoded the bytes, such as a long written as a string. Decode a sample row by hand before trusting a catalog.
- Authentication failures on a Kerberized cluster. Executors need HBase delegation tokens. Spark obtains them when HBase is on the classpath and
spark.security.credentials.hbase.enabledis true (the default); long jobs also need keytab-based renewal. - Bulk load left half done. HFiles were written but the load step failed, so the staging directory holds files that are not in the table. The load can be retried on the same directory; delete it only after success.
Trade-offs
| Approach | Best for | Cost |
|---|---|---|
| DataFrame source | SQL reads with row-key ranges; moderate writes | Selectivity depends on row key; writes are RPC Puts |
| HBaseContext bulkGet / bulkPut | Enrichment lookups, keyed updates from Spark | More code; batch sizes need tuning |
| Bulk load (HFiles) | Backfills, nightly rebuilds, very large writes | Shuffle and sort in Spark; bypasses WAL and WAL-based replication |
| Apache Phoenix | SQL with secondary indexes and its own type encoding | Another layer to run; tables must be Phoenix-managed |
If you need secondary indexes or SQL from many clients, look at Apache Phoenix. For how Spark data sources negotiate pushdown in general, see Spark Data Source V2, and for keeping scan tasks close to their data, HBase region locality. The scan settings above behave as described in HBase scans.
What to do next
- Build
hbase-sparkfor your exact Spark, Scala and HBase versions and puthbase-site.xmlin$SPARK_CONF_DIR. - Decide on column-filter pushdown: deploy the three jars to the RegionServers, or disable it explicitly.
- Write a catalog for one table, read ten rows, and decode them by hand to confirm the types.
- Run your most common query with
explain()and confirm row-key predicates become scan bounds. - Create and pre-split any table Spark writes to with split points from your own key distribution, not
newTable's default range, and salt sequential keys. - Move backfills and full rebuilds to
hbaseBulkLoadorhbaseBulkLoadThinRows, and keep the staging directory until the load succeeds. - Cap executors for jobs that scan live tables, and watch RegionServer latency while they run.