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.
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.
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.
Sizing each round trip
Every next RPC returns a batch of results, and four settings decide where a batch ends.
| Setting | Unit | Default (HBase 2.x) | What it controls |
|---|---|---|---|
setCaching(n) / hbase.client.scanner.caching | Rows per RPC | 2147483647 (no row cap) | Upper bound on rows returned by one RPC |
setMaxResultSize(b) / hbase.client.scanner.max.result.size | Bytes per RPC | 2 MiB client, 100 MiB server cap | Upper bound on bytes; in practice the real limit |
setBatch(n) | Cells per Result | Unlimited (whole row) | Splits wide rows into several Result objects |
setLimit(n) | Rows in total | None | Ends 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 readsetStartStopRowForPrefixScan 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
nextcalls. Process faster, shrink batches, or renew the lease. - Client out-of-memory on wide rows. Whole rows are being returned. Use
setBatchor 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
- Enable scan metrics on your three most common scans and record RPC calls, bytes, and rows scanned versus returned.
- Remove any hard-coded
setCachingfrom old code, then tunesetMaxResultSizeand addsetLimitwhere you need only the first rows. - Name only the families and columns each scan needs.
- Mark batch and export scans with
setCacheBlocks(false)and stream reads, and cap their parallelism. - Check that slow per-row processing cannot exceed the 60-second lease; move it off the scan loop if it can.
- Revisit row-key design for any query whose rows scanned greatly exceed rows returned; see RegionServer architecture for how regions serve those reads.