HBase is good at random reads and writes by row key. Spark is good at moving a lot of data through a cluster of executors. Teams want both in one job: enrich a Spark DataFrame with profile rows that live in HBase, write the output of a nightly aggregation back as serving rows, or scan a table to train a model. The Spark HBase connector is the piece that makes this practical without each job reinventing connection handling, region-aware partitioning and bulk loading.

This article covers the connector maintained by the Apache HBase project in the hbase-connectors repository, module hbase-spark. It has two layers: an RDD API built around HBaseContext, and a DataFrame source registered as org.apache.hadoop.hbase.spark. You may also meet the older Hortonworks connector (SHC, shc-core) in legacy code; it used a JSON catalog and a different format string, and new work should not start on it. By the end you should know how each layer reaches HBase, how to tell whether a filter is really pushed down, which write path fits which job, and what breaks in production.

How the connector reaches HBase

Start with connections, because most connector problems are connection problems. An HBase Connection is heavyweight: it holds a ZooKeeper session or registry client, a cache of region locations from hbase:meta, and thread pools. Creating one per record, or even per task, floods ZooKeeper and the meta region. HBaseContext solves this. You build it once on the driver from a Hadoop Configuration; it broadcasts that configuration (and, on secure clusters, the delegation credentials) to the executors, and each executor keeps one cached connection that every task on that executor reuses.

Reads follow the table's layout. When you scan through hbaseRDD, the input format creates one Spark partition per region, so a table with 240 regions gives you 240 tasks, each scanning one key range against one RegionServer. The DataFrame source does the same, but first removes regions whose key range cannot match your row-key predicates. Writes go the other way: each task opens a buffered mutator on its executor's connection and sends batches of puts to whichever RegionServers own those keys. The driver coordinates; it never carries the data.

Spark HBase connector: who holds the connection, who talks to which regionSpark driverHBaseContext(sc, conf)Broadcast confighbase-site.xml + credentialsExecutor 1one cached ConnectionExecutor 2one cached ConnectionExecutor None cached ConnectionRegionServer Aregions [a, h)RegionServer Bregions [h, p)RegionServer Cregions [p, z)ZooKeeper / hbase:metaregion locationsscan/putscan/putscan/putlookupRead pathone Spark partition per region (or pruned range)Write pathputs via BufferedMutator, or HFiles + bulk loadThe driver never touches data. Executors open one connection each and talk directly to the RegionServers that own the rows.
HBaseContext broadcasts configuration; each executor reuses one connection and talks to RegionServers directly.

Getting it onto the classpath

The connector has to match three version lines: Spark (and its Scala binary version), HBase, and Hadoop. Check Maven Central for an hbase-spark artifact built for your Spark and Scala versions before you plan on one. Many teams build from source with the flags the project README documents, then ship the jars with the job:

mvn -Dspark.version=3.3.1 -Dscala.version=2.12.15 -Dscala.binary.version=2.12 \
    -Dhbase.version=2.4.15 -Dhadoop-three.version=3.3.2 -DskipTests clean package

spark-submit --jars hbase-spark.jar,hbase-spark-protocol-shaded.jar \
  --packages org.apache.hbase:hbase-shaded-mapreduce:2.4.15 \
  --files /etc/hbase/conf/hbase-site.xml  my-job.jar

Two configuration rules matter. First, the client side needs hbase-site.xml (at minimum the ZooKeeper quorum or registry settings) on the driver and executor classpath, or passed into the Configuration you give HBaseContext. Second, column-filter pushdown runs code inside the RegionServers. The README states that scala-library, hbase-spark and hbase-spark-protocol-shaded must be on the RegionServer classpath for it to work. If you cannot change the servers, turn it off with .option("hbase.spark.pushdown.columnfilter", false). Without either step, the first filtered query fails with NoClassDefFoundError for a connector class, thrown from the server.

The RDD API: HBaseContext

The RDD layer is the one to reach for when your logic is per row and you want full control of the HBase objects. The main methods on HBaseContext, as listed in the Scaladoc, are:

  • bulkPut[T](rdd, tableName, f: T => Put): write every element as a Put.
  • bulkDelete[T](rdd, tableName, f: T => Delete, batchSize): delete in batches.
  • bulkGet[T, U](tableName, batchSize, rdd, makeGet, convertResult): look up many keys and return an RDD of results.
  • hbaseRDD(tableName, scan): scan a table into RDD[(ImmutableBytesWritable, Result)].
  • foreachPartition and mapPartitions: run your own code with the executor's Connection passed in.
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName}
import org.apache.hadoop.hbase.client.{Get, Put, Result, Scan}
import org.apache.hadoop.hbase.spark.HBaseContext
import org.apache.hadoop.hbase.util.Bytes

val conf = HBaseConfiguration.create()          // reads hbase-site.xml from the classpath
val hc   = new HBaseContext(spark.sparkContext, conf)
val profiles = TableName.valueOf("crm:profiles")
val cf = Bytes.toBytes("p")

// 1. Write: one Put per (userId, tier) pair.
val tiers = spark.sparkContext.parallelize(Seq(("u001", "gold"), ("u002", "silver")))
hc.bulkPut[(String, String)](tiers, profiles, { case (id, tier) =>
  new Put(Bytes.toBytes(id)).addColumn(cf, Bytes.toBytes("tier"), Bytes.toBytes(tier))
})

// 2. Point lookups: batches of 100 Gets per round trip.
val ids = spark.sparkContext.parallelize(Seq("u001", "u002", "u999"))
val found = hc.bulkGet[String, Option[String]](profiles, 100, ids,
  id => new Get(Bytes.toBytes(id)).addColumn(cf, Bytes.toBytes("tier")),
  r  => if (r.isEmpty) None else Some(Bytes.toString(r.getValue(cf, Bytes.toBytes("tier")))))

// 3. Scan one column over a key range: one partition per region in the range.
val scan = new Scan().withStartRow(Bytes.toBytes("u0")).withStopRow(Bytes.toBytes("u1"))
  .addColumn(cf, Bytes.toBytes("tier")).setCaching(500).setCacheBlocks(false)
val rows = hc.hbaseRDD(profiles, scan).map { case (k, r) => Bytes.toString(k.get()) }

Two settings in that scan are worth copying. setCaching(500) fetches 500 rows per RPC instead of the client default, which cuts round trips on a large scan. setCacheBlocks(false) stops a one-off full scan from evicting the hot blocks that online readers depend on. The HBase scans deep dive explains both settings at the RegionServer level.

The DataFrame source and column mapping

The DataFrame source lets you query HBase with Spark SQL. You describe the mapping from DataFrame columns to HBase cells in one string, hbase.columns.mapping, with entries of the form name TYPE family:qualifier; the row key uses the special form :key.

val mapping = "user_id STRING :key, tier STRING p:tier, ltv DOUBLE p:ltv, last_seen LONG a:ts"

val profiles = spark.read.format("org.apache.hadoop.hbase.spark")
  .option("hbase.columns.mapping", mapping)
  .option("hbase.table", "profiles")
  .option("hbase.namespace", "crm")
  .option("hbase.spark.use.hbasecontext", false)   // build connections from conf, not a pre-made context
  .load()

profiles.createOrReplaceTempView("profiles")
spark.sql("SELECT user_id, ltv FROM profiles " +
  "WHERE user_id >= 'u0420' AND user_id < 'u0430' AND tier = 'gold'").show()

// Writing uses the same mapping; the target table and families must already exist.
scored.select("user_id", "tier", "ltv")
  .write.format("org.apache.hadoop.hbase.spark")
  .option("hbase.columns.mapping", "user_id STRING :key, tier STRING p:tier, ltv DOUBLE p:ltv")
  .option("hbase.table", "profiles").option("hbase.namespace", "crm")
  .option("hbase.spark.use.hbasecontext", false)
  .save()

The types in the mapping decide how values become bytes, and HBase compares bytes, not numbers. The connector's default encoding uses HBase's Bytes.toBytes forms. A STRING key sorts the way you expect for fixed-width ids. An INT or LONG key is stored big-endian in two's complement, so in HBase's unsigned byte order negative numbers sort after positive ones. Any range you build yourself in a raw Scan, and any other system reading the table, sees that order, so a range crossing zero is a trap. Prefer non-negative or offset numeric keys. If another system writes the table, check its encoding too: a column written as a decimal string by a Java service and mapped as DOUBLE here will decode as garbage, not fail.

Pushdown: what really reaches the RegionServer

Pushdown is the reason to use the DataFrame source rather than scanning everything and filtering in Spark. It happens at two levels, and it is worth knowing which one your query hits.

Row-key predicates become key ranges. In the query above, user_id >= 'u0420' AND user_id < 'u0430' turns into a scan from u0420 to u0430. Regions outside that range are not scanned at all, so the job runs one or two tasks instead of 240. Equality on the key becomes a Get. This works on the client and needs nothing on the servers.

Column predicates become a server-side filter. tier = 'gold' cannot narrow the key range, so the connector sends a filter that RegionServers run against each row, returning only matches. That saves network and Spark CPU, but the RegionServers still read every row in the range. This is the part that needs the connector jars on the servers.

A worked example shows the difference. The table holds 200 million profiles across 240 regions of about 830,000 rows each, and 4 percent are gold. A query filtering only on tier = 'gold' still reads 200 million rows from disk on the RegionServers; pushdown cuts the 200 million rows shipped to Spark to 8 million. A query on a key range covering 10 regions reads about 8.3 million rows, which is 24 times less I/O, whether or not the tier filter is pushed. The lesson: design the row key so your common queries are range or prefix queries. Column filters save bandwidth, not disk reads. Run explain() and check the scan's PushedFilters, and compare the RegionServer read-request metrics while the job runs; those numbers show what was really scanned. The other read options, such as caching, batch size and time range, live in HBaseSparkConf (for example QUERY_CACHEDROWS, QUERY_BATCHSIZE, TIMERANGE_START and MAX_VERSIONS). Look up their string keys in the Scaladoc for the version you build rather than copying keys from a blog post.

Writing: puts versus bulk load

There are two write paths, and they suit different jobs.

Puts (bulkPut or a DataFrame save()) go through the normal write path: WAL, MemStore, flush. They are visible immediately, replicate to peer clusters, and fire coprocessors. The cost is load on the RegionServers. A job writing 50 million rows from 200 tasks at full speed can push MemStores into flush storms and blocking, and stall online writers. Cap the job's concurrency (fewer partitions, or coalesce before the write), and watch flush queue length and the count of blocked updates while it runs. HBase write performance lists the server settings involved.

Bulk load (bulkLoad or bulkLoadThinRows) writes HFiles directly from Spark, sorted and partitioned to match the table's regions, and then hands those files to HBase. It skips the WAL and MemStore, so it is far cheaper for the cluster, but it does not replicate through normal replication and it bypasses the write-path coprocessor hooks. The two methods differ in shape. bulkLoad takes a function that emits one (KeyFamilyQualifier, value) pair per cell, and it includes the qualifier in the shuffle sort, which handles very wide rows. bulkLoadThinRows takes one record per row (FamiliesQualifiersValues), and the project describes it as faster for rows with fewer than about 1,000 columns. Both only write the files; you then run the load step yourself (LoadIncrementalHFiles in older releases, BulkLoadHFiles in HBase 2.x). HBase bulk load covers what that step guarantees and what happens when a region splits in between.

import org.apache.hadoop.hbase.spark.{FamilyHFileWriteOptions, KeyFamilyQualifier}
import org.apache.hadoop.hbase.HConstants

val staging = "hdfs:///tmp/hfiles/profiles_20261003"
hc.bulkLoad[(String, String, Double)](rowsRdd, TableName.valueOf("crm:profiles"),
  r => Iterator(
    (new KeyFamilyQualifier(Bytes.toBytes(r._1), cf, Bytes.toBytes("tier")), Bytes.toBytes(r._2)),
    (new KeyFamilyQualifier(Bytes.toBytes(r._1), cf, Bytes.toBytes("ltv")),  Bytes.toBytes(r._3))),
  staging,
  new java.util.HashMap[Array[Byte], FamilyHFileWriteOptions](),  // per-family compression etc.
  false,                                                           // compactionExclude
  HConstants.DEFAULT_MAX_FILE_SIZE)                                // roll HFiles at this size
// then: BulkLoadHFiles.create(conf).bulkLoad(TableName.valueOf("crm:profiles"), new Path(staging))

When not to use the connector

The connector is not always the right read path. If a job reads a whole table to train a model or build a report, scanning live RegionServers competes with online traffic for disk, CPU and block cache. Two alternatives are often better. Take a snapshot and read it with TableSnapshotInputFormat through newAPIHadoopRDD, which reads HFiles straight from HDFS or object storage without involving the RegionServers. Or export the table once to Parquet and let every analytic job read that. Use the live connector when you need current data, key-range reads, or point lookups. If you need SQL with secondary indexes, Apache Phoenix is a separate layer with its own Spark integration. Mixing raw connector writes with Phoenix-managed tables corrupts the index, so pick one owner per table.

Failure modes

  • Connection storms. Code that calls ConnectionFactory.createConnection inside map opens one connection per task or per row, ZooKeeper connections hit their per-host limit, and the job fails with connection-loss errors. Use HBaseContext or foreachPartition.
  • Scanner timeouts. A task that does slow work between next() calls lets the scanner lease expire, and the next call fails with ScannerTimeoutException or UnknownScannerException. Lower caching, move slow work after the scan, or raise the scanner timeout.
  • Region moves mid-job. A split or balancer move during a scan causes NotServingRegionException; the client retries and usually recovers, but many at once mean you should pause the balancer for big jobs.
  • Pushdown class errors. NoClassDefFoundError from the server means the connector jars are missing on the RegionServers; disable column-filter pushdown or deploy them.
  • Silent type mismatches. Mapping types that differ from what writers used return wrong values, not errors. Validate a sample against a known row.
  • Kerberos. On secure clusters, executors need delegation tokens. Long jobs need a keytab and principal passed to spark-submit so tokens can be renewed.

What to do next

  1. Confirm which connector your code uses (Apache hbase-spark or legacy SHC) and that its build matches your Spark, Scala and HBase versions.
  2. Grep your jobs for connections created inside map or per task, and replace them with HBaseContext or foreachPartition.
  3. For each DataFrame read, run explain() and confirm the row-key predicate became a range; check RegionServer read-request counts during a test run.
  4. Decide on column-filter pushdown: deploy the server-side jars or set hbase.spark.pushdown.columnfilter to false, and record the decision.
  5. Classify every write job as puts or bulk load, cap the concurrency of put jobs, and alert on flush queue length and blocked updates.
  6. Move full-table analytic reads to snapshots or a Parquet export.
Key takeaway: The Spark HBase connector gives each executor one shared connection and maps Spark partitions to HBase regions, through either an RDD API built on HBaseContext or a DataFrame source configured by hbase.columns.mapping. Row-key predicates prune regions and save disk reads; column filters only save bandwidth and need connector jars on the RegionServers. Use puts for small, visible writes and bulk load for large ones, and read whole tables from snapshots rather than live servers.