Salting is the standard fix for an HBase table whose row keys arrive in order: timestamps, sequence numbers, ever-increasing IDs. HBase stores rows sorted by key and splits the key space into contiguous regions, so ordered keys all land in the last region, on one RegionServer, while the rest of the cluster idles. A salt is a short prefix computed from the key that deals consecutive rows into separate buckets. The idea fits in one sentence; the details are where teams get hurt. Pick the wrong salt input and Gets can no longer find rows. Pick the wrong bucket count and you are stuck with it, because changing it means rewriting the table. And every range scan quietly becomes N scans.
This article is the implementation manual; the concept and other fixes are covered in HBase hotspotting. Here we build a salted table end to end, from the determinism rule to migration.
The one rule: the salt is a function of the key
A salted row key is salt + naturalKey, where salt = hash(input) mod N and N is the number of buckets. The salt is normally one byte, which allows up to 256 buckets. The rule that makes salting work is that the salt must be deterministic: anyone who knows the natural key can recompute the salt and so rebuild the full row key.
Teams sometimes reach for a random prefix, or a round-robin counter, because it spreads writes perfectly. It also makes the row unfindable. A Get for device a91f at time T now has to try all N prefixes, which turns every point read into N reads. Worse, if the same logical record is written twice with different random prefixes, you get two copies, and the version semantics HBase gives you per row are gone. A deterministic salt keeps one natural key mapped to exactly one row key, so overwrites, deletes, increments and check-and-put all behave as they would on an unsalted table.
Salting is also not the same as replacing the key with its hash. Hashing the whole key spreads writes just as well but destroys ordering entirely: no range scan is possible at all. A salt keeps the natural key intact after the prefix, so each bucket is still a sorted run that you can scan.
Do you need a salt at all?
Salting fixes one specific shape: a leading key component that increases over time. If your key already starts with a high-cardinality, evenly distributed field, such as a user ID that is itself a UUID or hash, writes are already spread and a salt only adds read cost.
Salting earns its place when queries genuinely need time first in the key: "everything in the last five minutes, across all devices" is a range only if time leads.
What to hash
The salt input decides which reads can avoid the fan-out. Two choices cover nearly every design.
| Salt input | Write spread | Reads that hit one bucket | Risk |
|---|---|---|---|
Whole natural key (ts + deviceId) | Best: even a single busy device spreads | Gets only | Every scan fans out to N |
Entity only (deviceId) | Good if many entities write at once | Gets, and scans when the entity also leads the key after the salt | One hot entity is one hot bucket |
The hash itself must be stable across processes, JVM versions and client languages. java.util.Arrays.hashCode(byte[]) is specified by the Java documentation and stable, but a Python or Go writer will have to reimplement it exactly. Many teams choose MurmurHash3 or CRC32 with the variant written down, and pin test vectors, fixed keys with the salts the reference implementation computes, in every client's test suite. Never use a language's built-in string hash where it is randomised per process, as Python's hash() is.
Choosing the bucket count
Here is the arithmetic most guides skip. Inside each bucket the natural key still increases, so each bucket behaves like a small unsalted table: only its last region receives appends. With N buckets you therefore have exactly N concurrent write tails, however large the table grows and however many times its regions split. Splitting adds regions behind the tails; it never adds a tail.
That gives two bounds. The lower bound: N should be at least the number of RegionServers you want sharing the write load, and preferably a small multiple of it, so the balancer has room to place tails evenly. The upper bound: every range scan costs N scanner opens and N streams to merge, and a query that wants the newest K rows must pull up to K rows from each bucket. Small N means cheap scans and limited write parallelism; large N means the reverse.
| N | Write tails | Scanner opens per range query | Fits |
|---|---|---|---|
| 4 | 4 | 4 | Small cluster, scan-heavy reads |
| 16 | 16 | 16 | 10-15 RegionServers, mixed load |
| 64 | 64 | 64 | Large cluster, writes dominate, few scans |
A workable default is one to two times the current RegionServer count, rounded to a power of two. Because N cannot change cheaply, size it for the cluster you expect in two years, not the one you have today.
Pre-splitting on salt boundaries
A new table has one region. Unless you pre-split, all N buckets start inside it and you have the hotspot again until splits catch up. Create the table with N regions whose start keys are the salt bytes 1 to N-1. HBase compares row keys as unsigned bytes, so (byte) 200 sorts after (byte) 100 even though Java's signed byte would call it negative.
TableDescriptor desc = TableDescriptorBuilder.newBuilder(TableName.valueOf("telemetry"))
.setColumnFamily(ColumnFamilyDescriptorBuilder.newBuilder(Bytes.toBytes("d"))
.setCompressionType(Compression.Algorithm.ZSTD)
.build())
.build();
int buckets = 16;
byte[][] splits = new byte[buckets - 1][];
for (int i = 1; i < buckets; i++) {
splits[i - 1] = new byte[] { (byte) i }; // region i starts at salt byte i
}
try (Connection conn = ConnectionFactory.createConnection(conf);
Admin admin = conn.getAdmin()) {
admin.createTable(desc, splits);
}After creation, check in the HBase UI or with list_regions 'telemetry' in the shell that the N regions sit on different RegionServers. Regions will split further as each bucket grows; see regions and splits for the split policies that decide when.
Writing: one key builder, shared by everyone
Put the key logic in one small class that every writer and reader uses. Fixed-width fields keep the sort order meaningful; a variable-length field must be last or delimited.
public final class TelemetryKey {
public static final int BUCKETS = 16; // fixed at table creation
private static final int DEV_LEN = 16;
/** Row key = salt(1) | epochMillis(8, big-endian) | deviceId(16). */
public static byte[] rowKey(long epochMillis, byte[] deviceId) {
if (deviceId.length != DEV_LEN) throw new IllegalArgumentException("deviceId");
byte[] natural = Bytes.add(Bytes.toBytes(epochMillis), deviceId);
return Bytes.add(new byte[] { salt(natural) }, natural);
}
/** Salt over the whole natural key, so even one busy device spreads. */
public static byte salt(byte[] natural) {
return (byte) Math.floorMod(Arrays.hashCode(natural), BUCKETS);
}
}
// Writer
Put put = new Put(TelemetryKey.rowKey(event.ts(), event.deviceId()));
put.addColumn(D, Bytes.toBytes("v"), event.payload());
mutator.mutate(put); // BufferedMutator batches across regionsMath.floorMod matters: % on a negative hash returns a negative number, which would produce salt bytes outside 0 to N-1 and rows that land in the wrong region. Epoch milliseconds are positive, so Bytes.toBytes(long) sorts correctly; signed values would need their sign bit flipped first.
Reading: Gets stay cheap, scans fan out and merge
A Get rebuilds the key and costs one RPC, exactly as on an unsalted table. A time-range scan must run once per bucket. If the caller only needs the rows, not their order, issue the N scans in parallel and concatenate. If order matters, merge them with a priority queue keyed on the row key with the salt byte skipped.
public static List<Result> scanRange(Table table, long from, long to, int limit) throws IOException {
List<ResultScanner> scanners = new ArrayList<>();
PriorityQueue<Object[]> heads = new PriorityQueue<>((a, b) -> {
byte[] x = ((Result) a[0]).getRow(), y = ((Result) b[0]).getRow();
return Bytes.compareTo(x, 1, x.length - 1, y, 1, y.length - 1); // ignore salt byte
});
try {
for (int b = 0; b < TelemetryKey.BUCKETS; b++) {
byte[] salt = { (byte) b };
Scan scan = new Scan()
.withStartRow(Bytes.add(salt, Bytes.toBytes(from)))
.withStopRow(Bytes.add(salt, Bytes.toBytes(to))) // exclusive
.setCaching(Math.min(limit, 1000))
.setLimit(limit); // no bucket can contribute more than limit rows
ResultScanner s = table.getScanner(scan);
scanners.add(s);
Result r = s.next();
if (r != null) heads.add(new Object[] { r, s });
}
List<Result> out = new ArrayList<>();
while (!heads.isEmpty() && out.size() < limit) {
Object[] h = heads.poll();
out.add((Result) h[0]);
Result next = ((ResultScanner) h[1]).next();
if (next != null) heads.add(new Object[] { next, h[1] });
}
return out;
} finally {
scanners.forEach(ResultScanner::close);
}
}In a latency-sensitive service, open the scanners from a thread pool or use AsyncTable so the N round trips overlap. The per-bucket setLimit bounds the work. Scanner caching is explained in HBase scans.
Key layout decides which queries stay cheap
Salting only redistributes rows; the fields after the salt still decide which queries are ranges. Write down your top queries and test each layout against them before creating the table.
| Layout | Global time window | One device's history | Write hotspot? |
|---|---|---|---|
deviceId | ts (no salt) | Full table scan | One range scan | No, if device IDs are spread |
salt(ts,dev) | ts | dev | N range scans | N scans plus a device filter | No |
salt(dev) | dev | ts | Full table scan | One range scan in one bucket | Only for one very hot device |
When both query shapes matter at volume, the common answer is two tables written together: a salted time-first table for windows and an entity-first table for histories.
Worked example: an IoT telemetry table
A fleet of 200,000 devices reports every few seconds into an unsalted table keyed ts | deviceId. The cluster has 12 RegionServers, yet write latency climbs and one server's write request count dwarfs the others: every new row sorts after every old one, so one region takes all appends, splits, and the new daughter takes them all again.
The team's main query is "all readings in a five-minute window" for anomaly detection, so time must lead. They choose N = 16: more than the 12 servers, a power of two, with headroom to grow to around 16 servers without leaving one idle. They hash the whole natural key so a chatty device cannot pin a bucket, pre-split into 16 regions, and backfill with the method in the next section.
Afterwards writes spread over 16 tail regions on all 12 servers. The window query opens 16 small scanners instead of one, and device-history lookups move to a second table keyed deviceId | ts.
Changing the bucket count later
Every row's salt depends on N, so a new N moves almost every row to a different key. There is no in-place change; it is a migration.
- Create the new table, pre-split for the new N.
- Change writers to dual-write to both tables, with the new key builder for the new table.
- Backfill history with a MapReduce or Spark job that reads a snapshot of the old table and rewrites each row's key. Writing HFiles and bulk loading them is far cheaper than Puts for this; see HBase bulk load. The HFiles must be partitioned by the new table's region boundaries.
- Compare row counts per time range between the tables, then switch readers.
- Stop dual writes and drop the old table after a safety period.
Failure modes
- Reader and writer disagree on the salt. A Python ingest job and a Java reader using different hashes produce rows that Gets never find, and nothing errors. Shared test vectors catch it.
- Forgetting the fan-out. A new endpoint scans only bucket 0, or only uses
withStartRowwith no salt byte, and silently returns about 1/N of the data. - No pre-split. The table starts with one region, so the hotspot is back for the first hours or days, exactly when a launch generates the most traffic.
- Tails clustered on one server. If several bucket tails end up on the same RegionServer, that server is still hot. Check region placement after splits and rebalance if needed.
Letting Phoenix do it
If you query HBase through Apache Phoenix, CREATE TABLE ... SALT_BUCKETS = 16 adds a hidden salt byte, pre-splits the table and runs the fan-out and merge inside its parallel scans for you. Phoenix accepts values from 1 to 256. The same rules apply: the bucket count is fixed at creation, and non-Phoenix clients must not write to the table, because they will not compute Phoenix's salt. Apache Phoenix on HBase covers how its planner uses the salt.
What to do next
- Look at the leading bytes of real row keys and the per-region write request counts; salt only if a monotonic field leads and one region takes the writes.
- List your top queries and choose the key layout that keeps them ranges; consider a second table for the other access path.
- Pick the salt input and a stable hash, and pin test vectors in every client language.
- Choose N as one to two times the RegionServer count you expect, rounded to a power of two.
- Create the table pre-split on salt bytes 1 to N-1 and confirm the regions are spread.
- Wrap key building, Gets and the merging scanner in one shared library with per-bucket limits.
- Write down the migration plan for a new N before you need it.