Ad tech was one of the workloads HBase was made for: billions of anonymous device identifiers, a few kilobytes of state per device, a very high read rate on a strict deadline, and a write stream that is mostly counters and append-only events. It is also a workload where the obvious schema fails in specific, predictable ways: one campaign's budget becomes a hot row, retried increments double-count impressions, and a privacy deletion leaves data sitting in HFiles and snapshots long after the API returned success.

This article builds the HBase side of a small demand-side platform from first principles: the four access patterns an ad stack puts on the store, a table design and client code for each, sizing numbers for a mid-sized deployment, and the failure modes to test before peak traffic finds them. The examples use the standard HBase 2.x Java client.

What an ad stack asks of the store

A real-time bidding exchange sends the bidder a request describing an impression opportunity, and the request carries a deadline. In OpenRTB that deadline is the tmax field, set per request by the exchange, and it covers the whole round trip including the network. What is left after network time and auction logic is the store's budget, measured at the tail. A late bid is not a slow bid; it is a lost one.

Four patterns hit the store, and they have very different shapes:

PatternOperationLatency needConsistency need
Profile lookupSingle-row Get by device keyStrict, on the bid pathStale by seconds is fine
Frequency cappingRead at bid time, Increment at impressionRead strict, write relaxedApproximately right, never wildly over
Budget pacingSpend counters per campaignRelaxedBounded overspend, not exact
AttributionPut impression, look up on clickRelaxedMust not lose impressions

The design rule that follows is simple to state: the bid path does reads only, and nothing on the bid path waits for a write. Every write is produced after the auction by the event pipeline, and the bidder learns about it on its next read. Accepting that one-way staleness is what makes the latency budget achievable at all.

Architecture: three tables, two paths

HBase in an ad stack: one cluster, four access patterns with different budgetsExchangebid request, tmaxBidderauction logic in memoryprofiles (read replicas)Get: segments + caps, one rowGetImpression / click beaconsfired after the auctionEvent pipelineKafka + stream jobprofiles: fc familyIncrement with nonceimpressions (salted)Put, TTL = attribution windowPacing servicein-memory spend, sharded rowsbudget (sharded)N counter rows per campaignspend eventsflush every few sBatch / MLsegments rebuilt nightlybulk load HFilesBid path: reads only, inside a hard deadline. Write path: asynchronous.
Figure 1. Bid-path reads go to the profiles table, served by read replicas. Impressions, clicks and spend arrive later through the event pipeline, and segments arrive in bulk from batch jobs.

Three tables are enough. profiles holds one row per device with two column families: seg for audience segments written by batch jobs, and fc for frequency-cap counters written by the stream; they differ in write rate and TTL. impressions holds one row per served impression for the attribution window, and budget holds sharded spend counters. Campaign configuration stays out of HBase: it is small and belongs in the bidder's memory.

Row keys and table layout

Device identifiers are already high-entropy if they are random advertising IDs, but many ID sources are not (sequential internal IDs, hashed emails with a common prefix after normalisation errors). Hash every incoming ID with a keyed hash before it becomes a row key. That does three jobs at once: it spreads rows evenly across regions, it means a raw identifier never sits in the store, and a deletion request can be served by hashing the same input to find the row.

// Row key: first 16 bytes of HMAC-SHA256(secret, normalised id). Fixed width, uniform, not reversible.
static byte[] profileKey(Mac hmac, String rawId) {
    byte[] h = hmac.doFinal(rawId.trim().toLowerCase(Locale.ROOT).getBytes(StandardCharsets.UTF_8));
    return Arrays.copyOf(h, 16);
}

// Table: profiles
//   family seg  : qualifier = segment id, value = score byte; VERSIONS 1; TTL 30 days; BLOOMFILTER ROW
//   family fc   : qualifier = campaignId|yyyymmdd, value = 8-byte counter; TTL 2 days; IN_MEMORY true
// Pre-split into regions on the hash prefix, so day one does not start with one hot region.

A uniform key makes pre-splitting easy: divide the 16-byte space into equal ranges. Impression IDs are usually UUIDs and also uniform; if yours are time-ordered, add a salt prefix as described in HBase salting. Never key impressions by timestamp, or every write lands on the last region.

Profile lookup on the bid path

The profile read is one Get that fetches both families for one row. Restrict it to the families you need, give it a hard timeout shorter than your bid budget, and treat a timeout as an empty profile rather than an error. An empty profile still lets the bidder bid on contextual signals; a thrown exception loses the auction.

Get get = new Get(profileKey(hmac, deviceId))
        .addFamily(SEG).addFamily(FC)
        .setConsistency(Consistency.TIMELINE);   // may be served by a secondary replica
Result r;
try {
    r = profiles.get(get);                       // Table built with a short operation timeout
} catch (IOException e) {
    metrics.profileTimeout.inc();
    r = Result.EMPTY_RESULT;                     // bid contextually, never fail the auction
}
if (r.isStale()) metrics.staleRead.inc();       // answered by a secondary, not the primary

Region replicas are what make tail latency survivable. With replication set on the table, each region has a primary and one or more read-only secondaries on other RegionServers. A TIMELINE Get goes to the primary first and, if it is slow to answer, also to the secondaries, returning the first answer flagged with isStale(). A garbage-collection pause or crash then costs freshness, not bids. The price is memory for the copies and secondaries that lag, which for caps means a few seconds of under-counting. The mechanics are covered in HBase region replicas.

The rest comes from the usual places: a ROW Bloom filter so a miss for an unknown device skips most HFiles, a block cache that holds the hot set, and few HFiles per region. Measure p99 from the bidder; the client sees retries, the server does not.

Frequency capping with increments

Frequency capping limits how often one device sees one campaign, for example three times a day. The read happens at bid time as part of the profile Get above: the fc family returns today's counters, and the bidder skips campaigns at their cap. The write happens later, when the impression beacon fires, as an atomic Increment on the row.

// Called by the stream job for every impression beacon.
byte[] q = Bytes.toBytes(campaignId + "|" + day);           // e.g. "c8812|20261004"
Increment inc = new Increment(profileKey(hmac, deviceId))
        .addColumn(FC, q, 1L);
inc.setTTL(Duration.ofDays(2).toMillis());                  // cell TTL, in milliseconds
profiles.increment(inc);                                    // client attaches a nonce; see below

Three details decide whether this is correct. First, Increment is not naturally idempotent: if the client times out after the server applied the increment and then retries, the count goes up twice. HBase addresses this with client nonces, enabled by default, which let the server recognise a retry of an operation it already applied for a bounded period. They do not cover your stream job replaying a beacon after a restart; that needs dedup on beacon ID upstream, or acceptance of a small over-count, the safe direction for capping.

Second, the cell TTL set with setTTL is in milliseconds, while a column family TTL is in seconds. Mixing them up produces counters that expire in seconds or live for decades. TTL and versions in HBase covers how expiry interacts with compaction.

Third, increments take a row lock. A normal device sees a few impressions an hour and the lock is free. A bot or a shared device ID (an unset identifier that defaults to all zeros is common) can receive thousands per second, and every increment on that row serialises. Filter known-invalid IDs before the store, and alarm on any row whose increment rate exceeds a human ceiling.

Budget pacing without a hot row

Budget pacing is where the naive design fails loudest. If every won impression increments one row per campaign, a large campaign makes that row the busiest object in the cluster, and since a row lives in one region, adding servers does not help. Two changes fix it.

  1. Aggregate in memory, flush periodically. Each pacing worker sums spend per campaign locally and flushes one increment every few seconds. Thousands of writes per second become one per worker per interval.
  2. Shard the counter. Store the campaign's spend across N rows, campaignId|00 to campaignId|15, with each worker writing to the shard chosen by its worker ID. Reading total spend is a short Scan over N rows, which the pacing service does on its own schedule, never on the bid path.
// Pacing worker: flush local sums every 5 s to a sharded counter row.
byte[] row = Bytes.toBytes(campaignId + "|" + String.format("%02d", workerId % SHARDS));
budget.increment(new Increment(row).addColumn(SPEND, today, microsSpent));

// Pacing service: read total spend with one short scan.
Scan s = new Scan().withStartRow(Bytes.toBytes(campaignId + "|"))
                   .withStopRow(Bytes.toBytes(campaignId + "|~"))
                   .addColumn(SPEND, today);
long total = 0;
try (ResultScanner rs = budget.getScanner(s)) {
    for (Result r : rs) total += Bytes.toLong(r.getValue(SPEND, today));
}

The shards share a prefix and sit in one region, which is fine once writes are aggregated. The cost is bounded overspend: spend still in worker memory is uncounted. With 5-second flushes and a campaign spending 10 currency units per second, up to about 50 units are unflushed, so pacing must stop at least that far below the cap.

Impressions and attribution

Attribution joins a click or conversion to the impression that caused it. The impression beacon writes one row keyed by impression ID with the campaign, creative, device key, price and timestamp; the click beacon carries the same impression ID, so the join is a single Get. Set the family TTL to the attribution window, for example seven days for clicks, so the table's size is bounded by traffic times window and never grows without limit.

Write impressions with the WAL on; skipping it per mutation saves little and lost impressions mean billing disputes. Load nightly segments into profiles:seg as bulk-loaded HFiles rather than Puts, so the refresh does not compete with live traffic for MemStore and WAL. HBase bulk load walks through the job.

Privacy deletion

A deletion request for a device must remove its profile, counters and impressions. HBase makes the API call easy and the guarantee hard. A Delete writes a tombstone; the old cells remain in HFiles until a major compaction rewrites the region, and they remain indefinitely in any snapshot or exported backup taken before the deletion. Also, impressions are keyed by impression ID, not device, so you cannot find them by device without a scan.

  • Keep a reverse index from device key to recent impression IDs, or let impressions age out by TTL and state that window in your privacy notice.
  • Track deletion as steps (tombstone, major compaction, oldest prior snapshot expired), not one API call.
  • Rotate the HMAC secret only with a migration plan: a new secret makes every existing row unreachable by ID.

Worked example: sizing a mid-sized bidder

Take a bidder that receives 300,000 bid requests per second at peak, of which half pass pre-filters and need a profile. That is 150,000 single-row Gets per second. With profiles averaging 1.5 KB and 2 billion devices, the table is about 3 TB before replication and compression. If 20 percent of devices produce 80 percent of requests, the hot set is around 600 GB, which spread over 30 RegionServers means 20 GB of block cache each, held in an off-heap bucket cache. About 5,000 Gets per second per server is modest; the tail is the hard part.

At a 15 percent win rate the cluster also takes about 22,000 impression Puts and as many cap increments per second, plus a few hundred aggregated budget writes. Load-test reads and writes together at peak, so you see whether the 45,000 writes per second lift the read p99.

Failure modes

FailureSymptomFix
Hot device rowOne region's latency spikes; increments queueDrop invalid IDs; rate-alarm per row
Hot campaign counterPacing lags; one RegionServer saturatedIn-memory aggregation plus sharded rows
Retried incrementsCaps hit early, spend over-reportedKeep nonces on; dedup beacons by ID upstream
TTL unit mix-upCounters vanish or never expireCell TTL in ms, family TTL in s; test both
Compaction stormRead p99 climbs during the daySchedule majors off-peak; watch store file count
GC pause on a primaryBid timeouts in burstsTIMELINE reads with replicas; off-heap cache
Deletion incompleteDeleted data found in a snapshotTrack compaction and snapshot expiry per request

Most of these show up first as hot regions. Diagnosing HBase hotspots has the metrics and the per-region request counts to watch.

Trade-offs

HBase earns its place when you already run Hadoop-family infrastructure, bulk load from batch jobs, and want consistent single-row operations such as increments. Its costs are operational: ZooKeeper, HDFS, compaction tuning and JVM heaps. Purpose-built key-value stores can give lower tails with less tuning, and managed wide-column services remove the operations work. Decide on your own traffic, not benchmarks.

What to do next

  1. Write down your bid budget from the exchange's tmax values and network time, and set the client operation timeout below it.
  2. Key profiles by a fixed-width keyed hash of the device ID and pre-split on it.
  3. Split segment and cap data into separate families with their own TTLs.
  4. Enable region replicas on the profiles table and read with TIMELINE consistency; graph stale-read rate.
  5. Move budget counting to in-memory aggregation with sharded counter rows.
  6. Confirm client nonces are enabled and add beacon dedup in the stream job.
  7. Load-test reads and writes together at peak, and record the bidder-side p99.
  8. Write the deletion runbook, including compaction and snapshot expiry, before the first request arrives.
Key takeaway: Keep the bid path read-only: one Get per request against a hashed device key, served by region replicas with TIMELINE consistency and a timeout that degrades to an empty profile. Write caps, impressions and spend asynchronously, keep increments safe with nonces and upstream dedup, and never let one campaign's budget live in one hot row. Treat privacy deletion as a tracked process that ends only when compaction and snapshot expiry have removed the data.