HBase stores rows sorted by row key and splits the key space into regions, each served by one RegionServer. A Get reads one row. A Scan reads a contiguous range of rows, in key order, and it is the operation behind every time-series read, every per-customer history page and every batch export. When HBase feels slow, a badly shaped scan is the most common cause.

A scan is not a single request. It is a conversation: the client opens a scanner on one region, pulls rows in a series of next calls, moves to the next region, and repeats. This article explains that conversation from the client loop down to the file reads, shows how to size each round trip, works out how many RPCs a real scan needs, and covers the timeouts, read modes and parallelisation patterns that decide whether a scan takes seconds or hours.

Advertisement

The shape of a scan

You describe a scan by a start row (inclusive by default), a stop row (exclusive by default), the column families or columns you want, optional filters, a time range and a version count. Because rows are sorted, HBase can jump to the start row and read forward until the stop row, touching nothing outside the range. That is the whole reason row-key design matters: a query whose rows are adjacent in key order is one cheap scan, while a query whose rows are scattered is either many Gets or a scan over far more data than it returns.

The client first finds which region holds the start row by consulting hbase:meta (and caching the answer; see how a row key finds its RegionServer). It opens a scanner on that region, reads until the region ends or the stop row is reached, closes it, and opens a scanner on the next region. Regions are visited one after another, so a plain scan uses one RegionServer at a time no matter how many the range spans.

One scan: the client loop, the region hop, and the server-side mergeClient ResultScannerlocal cache of ResultsRegionServer RPCopen / next / close, scanner idRegionScannerone per region, holds leasenext()rows or heartbeatcache empty: RPC; region ends: open next regionStoreScanner per familyKeyValueHeap mergesmallest next cell winsMemStoresegment scannersHFile 1pread or streamHFile 2pread or streamEach next() RPC stops atcaching rows, or maxResultSize bytes,or limit, or time budget (heartbeat)Rows arrive in key order within a region; regions are visited in key order (reverse order if reversed).
The client refills its local cache with next() RPCs; each RPC is bounded by rows, bytes, a row limit or a time budget. On the server, one StoreScanner per column family merges MemStore and HFile scanners through a heap so cells come out in key order.

What happens inside the RegionServer

On the server, a RegionScanner represents your scan on one region. Under it sits one StoreScanner per column family you asked for, and under each of those sit scanners over the family's MemStore and every HFile. Because each source is individually sorted, a heap merges them: whichever source has the smallest next cell supplies it, so the scan sees one sorted stream even though the data is spread across memory and several files. Newer versions of a cell shadow older ones, delete markers hide what they cover, and the time range and version limit are applied here.

Two consequences are worth remembering. First, the number of HFiles matters: every extra file is another source in the merge and another potential disk read, which is why read performance degrades between compactions. Second, asking for fewer column families means fewer StoreScanners, so a scan that names only the family it needs is cheaper than one that reads the whole row. Filters run inside this loop too, and the best ones tell the scanner to seek ahead instead of reading every cell; HBase filters in depth covers which ones do.

Each region scanner reads from a consistent point: writes that commit after the scanner opens are not visible to it. That consistency is per region, not across the whole scan, because each region's scanner opens at a different moment.

Advertisement

Sizing each round trip

Every next RPC returns a batch of results, and four settings decide where a batch ends.

SettingUnitDefault (HBase 2.x)What it controls
setCaching(n) / hbase.client.scanner.cachingRows per RPC2147483647 (no row cap)Upper bound on rows returned by one RPC
setMaxResultSize(b) / hbase.client.scanner.max.result.sizeBytes per RPC2 MiB client, 100 MiB server capUpper bound on bytes; in practice the real limit
setBatch(n)Cells per ResultUnlimited (whole row)Splits wide rows into several Result objects
setLimit(n)Rows in totalNoneEnds the whole scan after n rows

Because caching defaults to effectively infinite rows, the byte limit decides batch size in modern HBase: each RPC returns roughly 2 MiB. That is a sensible default. Old advice to 'set caching to 100' dates from releases where the default was one row per RPC and can now make scans slower by cutting batches short.

Wide rows need care. A single row with millions of cells can exceed any byte limit, and by default HBase returns a whole row together, so one row can dominate a response and the client's memory. setBatch caps the cells per Result, so one row may arrive as several Results. setAllowPartialResults(true) lets the server return a partial row when it hits the size limit, and the client sees Result.mayHaveMoreCellsInRow() set to true. Either way, your code must be ready to stitch a row back together, and you give up row-level atomicity for that row across results.

Use setLimit whenever you only need the first N rows. Without it, the server fills a 2 MiB batch even when you stop reading after ten rows.

Worked example: counting RPCs

A table holds a month of events for one customer: 5 million rows of about 400 bytes each, 2 GB in total, spread across 4 regions. With default settings each RPC carries about 2 MiB, or roughly 5,200 rows, so the scan needs about 960 next calls plus an open and close per region. If each call spends 2 ms in the network and 10 ms reading on the server, the scan takes around 12 seconds, and the client spends most of that time waiting.

Raise setMaxResultSize to 8 MiB and the call count drops to about 240; server time per call grows, but the per-call network overhead falls fourfold, so the scan finishes about 1.5 seconds sooner at the cost of 8 MiB batches in client memory. Enable setAsyncPrefetch(true) and the client fetches the next batch while your code processes the current one, overlapping the two. Split the scan across its 4 regions and run them in parallel, and wall-clock time falls roughly fourfold, because four RegionServers now work at once. Ask only for the one column you need, and the bytes per row, and so every other number, shrink with it. Confirm each change with scan metrics rather than intuition.

Writing the scan

import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.client.*;
import org.apache.hadoop.hbase.client.metrics.ScanMetrics;
import org.apache.hadoop.hbase.util.Bytes;

Scan scan = new Scan()
    .withStartRow(Bytes.toBytes("cust42#2026-09-01"))          // inclusive
    .withStopRow(Bytes.toBytes("cust42#2026-10-01"), false)    // exclusive
    .addColumn(Bytes.toBytes("e"), Bytes.toBytes("amt"))
    .setMaxResultSize(4L * 1024 * 1024)  // bytes per RPC; no row cap set
    .setCacheBlocks(true)            // small, hot range: keep blocks cached
    .setScanMetricsEnabled(true);

try (Table table = connection.getTable(TableName.valueOf("events"));
     ResultScanner scanner = table.getScanner(scan)) {
  for (Result r : scanner) {
    process(r);                      // keep this fast: the lease is ticking
  }
  ScanMetrics m = scanner.getScanMetrics();
  System.out.println(m.getMetricsMap());   // RPC_CALLS, BYTES_IN_RESULTS, ...
}

ResultScanner is Closeable; always close it. An abandoned scanner keeps server resources until its lease expires. getMetricsMap() returns counters such as RPC_CALLS, REMOTE_RPC_CALLS, REGIONS_SCANNED, BYTES_IN_RESULTS and MILLIS_BETWEEN_NEXTS, and server-side counts of rows scanned and rows filtered. A large gap between rows scanned and rows returned tells you the key range is too wide and filters are doing work the key design should do.

Leases, timeouts and heartbeats

An open scanner holds state on the RegionServer, so the server gives it a lease: if the client does not call next within hbase.client.scanner.timeout.period (60 seconds by default), the server discards the scanner and the next call fails. The classic cause is slow processing between calls: a client that pulls a 2 MiB batch and then spends 90 seconds writing each row to another system loses its scanner. Fixes, in order of preference: process faster or hand rows to a separate worker, shrink the batch so each is processed within the lease, call ResultScanner.renewLease() during long pauses, or raise the timeout on both client and servers.

The opposite problem arises when a selective filter makes the server read many cells without finding any to return. Rather than let the RPC run until it times out, the server checks its time budget every 10,000 cells (hbase.cells.scanned.per.heartbeat.check) and, when the budget is spent, returns an empty heartbeat response so the client knows the scan is alive and asks again. Frequent heartbeats in a trace are a sign that the scan is reading far more than it returns. With setNeedCursorResult(true), heartbeats also carry the current row position, so an application can report progress or resume from a known key.

Read modes, block caching, reversed and prefix scans

Each HFile can be read with positional reads (pread), which fetch exactly the requested block and suit short scans, or with a streaming read, which reads ahead sequentially and suits long ones. By default a scan starts with pread and switches to stream once it has read more than hbase.storescanner.pread.max.bytes, four times the block size by default. If you know a scan is long, setReadType(Scan.ReadType.STREAM) skips the ramp; for known-short scans, PREAD avoids opening a stream at all.

Every block a scan reads normally enters the block cache. A full-table export would evict the hot working set that serves your latency-sensitive Gets, so batch scans should call setCacheBlocks(false). See block cache architecture for how the cache prioritises blocks.

// Forward: every row for one customer. The helper computes start and stop.
Scan all = new Scan().setStartStopRowForPrefixScan(Bytes.toBytes("cust42#"));

// Reversed: newest first, first 20 rows only. A reversed scan starts at the
// HIGH key, and the prefix helper ignores the reversed flag, so set the
// bounds yourself: '$' is the byte after '#', the first key past the prefix.
Scan recent = new Scan()
    .setReversed(true)
    .withStartRow(Bytes.toBytes("cust42$"), false)
    .withStopRow(Bytes.toBytes("cust42#"), true)
    .setLimit(20);        // stop after 20 rows; no wasted RPCs past the answer

// A full-table batch job: do not pollute the block cache, read sequentially.
Scan export = new Scan()
    .setCacheBlocks(false)
    .setReadType(Scan.ReadType.STREAM)
    .setMaxResultSize(8L * 1024 * 1024);   // bigger batches for a long read

setStartStopRowForPrefixScan converts a prefix into forward start and stop rows (it ignores the reversed flag, so reversed prefix scans set bounds by hand, high key first); it replaces the deprecated setRowPrefixFilter, whose name suggested a filter although it never was one. Do not combine it with explicit start or stop rows. Reversed scans read regions and rows from high to low, which is how 'newest first' queries work when the key ends in an ascending timestamp; they are somewhat slower than forward scans because HFile blocks are laid out for forward reading, so if newest-first is the dominant query, store a reversed timestamp in the key instead.

Parallel scans and batch processing

A single scan uses one region at a time. To use the whole cluster, split the range at region boundaries and scan the pieces concurrently:

// Split a large scan along region boundaries and run the pieces in parallel.
// max, minStop and isEmptyRange are small byte-array helpers (Bytes.compareTo).
List<Scan> splitByRegion(Connection conn, TableName tn, byte[] start, byte[] stop)
    throws IOException {
  List<Scan> out = new ArrayList<>();
  try (RegionLocator loc = conn.getRegionLocator(tn)) {
    Pair<byte[][], byte[][]> keys = loc.getStartEndKeys();
    for (int i = 0; i < keys.getFirst().length; i++) {
      byte[] rs = keys.getFirst()[i], re = keys.getSecond()[i];
      byte[] s = max(rs, start);                 // clip to requested range
      byte[] e = minStop(re, stop);              // empty = open-ended
      if (isEmptyRange(s, e)) continue;
      out.add(new Scan().withStartRow(s).withStopRow(e)
          .setCacheBlocks(false));
    }
  }
  return out;   // submit each to an executor with bounded concurrency
}

Bound the concurrency: one hundred simultaneous scans against a 20-node cluster will saturate handlers and hurt every other client. MapReduce's TableInputFormat and the Spark connector do the same split automatically, one task per region. For very large offline jobs, scanning a table snapshot reads HFiles directly from the filesystem and bypasses the RegionServers entirely, which protects online traffic at the cost of reading a point-in-time copy.

If clients reach HBase through the Thrift or REST gateway rather than the Java client, scanner state lives in the gateway, which changes how failures and load balancing behave; see the Thrift and REST gateway article.

Failure modes and what they mean

  • Scanner lease expired or unknown scanner errors. Too long between next calls. Process faster, shrink batches, or renew the lease.
  • Client out-of-memory on wide rows. Whole rows are being returned. Use setBatch or partial results, and select fewer columns.
  • Latency-sensitive reads slow down during batch jobs. Block cache pollution or handler saturation. Disable block caching on batch scans, cap parallelism, or scan a snapshot.
  • A 'short' query scans for minutes. The key design forces a wide range and filters discard most rows. Compare rows scanned with rows returned, then fix the key or add an index table.
  • One RegionServer pegged by a scan. The range sits in one hot region. Parallel splitting cannot help; salting or presplitting the key space can.

What to do next

  1. Enable scan metrics on your three most common scans and record RPC calls, bytes, and rows scanned versus returned.
  2. Remove any hard-coded setCaching from old code, then tune setMaxResultSize and add setLimit where you need only the first rows.
  3. Name only the families and columns each scan needs.
  4. Mark batch and export scans with setCacheBlocks(false) and stream reads, and cap their parallelism.
  5. Check that slow per-row processing cannot exceed the 60-second lease; move it off the scan loop if it can.
  6. Revisit row-key design for any query whose rows scanned greatly exceed rows returned; see RegionServer architecture for how regions serve those reads.
Key takeaway: An HBase scan is a sequence of next() RPCs against one region at a time, each served by a heap merge over the MemStore and every HFile of the families you requested. In HBase 2.x the byte limit, not the row count, sizes each RPC. Scan only the key range and columns you need, use setLimit for top-N queries, keep per-row processing inside the lease, keep batch scans out of the block cache, split large scans by region, and use scan metrics to prove each change helped.