hot partitions
Hot Partitions
Learn why uneven key and traffic distributions overload individual database partitions, how hot keys differ from storage skew, how low-cardinality and monotonically increasing keys create bottlenecks, and how hashing, bucketing, write sharding, caching, aggregation, isolation, throttling, and repartitioning can distribute load safely.
Introduction
Distributed databases divide data into partitions so storage and request processing can be spread across several nodes.
Adding more partitions does not automatically guarantee balanced work. If a large share of requests or data targets one key range, tenant, product, timestamp, status, or entity, the corresponding partition can become overloaded while other partitions remain underused.
This overloaded partition is called a hot partition.
Core idea: A distributed database scales effectively only when its data and traffic can be distributed. A partition key that concentrates requests on one partition turns a distributed cluster into a system limited by one narrow bottleneck.
Hot partitions are primarily an access-pattern and data-modeling problem. Adding nodes cannot fully solve the issue when the same logical key continues sending all work to one partition.
Prerequisites
| # | Prerequisite | Why It Is Needed |
|---|---|---|
| 1 | NoSQL data models | Key-value, document, and wide-column systems commonly distribute records using partition keys. |
| 2 | Access-pattern design | Partition keys must support both query locality and workload distribution. |
| 3 | Secondary indexes | A secondary index can develop a hotspot independently of the base table. |
| 4 | Denormalization | Query-specific records and replicated summaries can redistribute read traffic. |
| 5 | Capacity estimation | Average traffic is insufficient when one key receives a disproportionate share. |
| 6 | Caching and rate limiting | These controls can protect hot read and write paths but introduce consistency trade-offs. |
What Is a Partition?
A partition is a logical portion of a dataset assigned to a storage and processing unit.
Distributed database
+------------------+
| Partition A |
| Keys 1 to 1,000 |
+------------------+
+------------------+
| Partition B |
| Keys 1,001-2,000 |
+------------------+
+------------------+
| Partition C |
| Keys 2,001-3,000 |
+------------------+
A partition can be mapped to a node, shard, tablet, split, replica group, or another database-specific unit.
Partitioning aims to distribute:
- Stored data
- Read requests
- Write requests
- CPU usage
- Memory consumption
- Disk input and output
- Network traffic
- Background maintenance
Partition Key
A partition key is the attribute or key combination used to determine where a record belongs.
Record:
tenantId = 17
courseId = 42
learnerId = 1042
Candidate partition keys:
tenantId
courseId
learnerId
tenantId + courseId
tenantId + learnerId
The chosen key affects:
- Data placement
- Request routing
- Query locality
- Maximum partition size
- Traffic distribution
- Range-scan behaviour
- Repartitioning complexity
What Is a Hot Partition?
A hot partition receives a disproportionate share of data, reads, writes, or processing work compared with other partitions.
Partition A:
10% CPU
100 requests per second
Partition B:
12% CPU
120 requests per second
Partition C:
98% CPU
8,000 requests per second
Result:
Partition C is hot.
The hot partition can reach its resource or throughput limit even though the complete cluster still has unused capacity.
Hot Partition vs Hot Key
| Concept | Meaning | Example |
|---|---|---|
| Hot key | One logical key receives a disproportionate number of requests | One popular course receives most page views |
| Hot partition | One partition receives disproportionate data or traffic | Several popular keys map to the same physical partition |
| Data skew | Stored data is distributed unevenly | One tenant owns most records |
| Traffic skew | Requests are distributed unevenly | One small but popular record receives most reads |
A hot key commonly creates a hot partition, but a partition can also become hot because several unrelated high-traffic keys happen to map to it.
Data Skew vs Traffic Skew
Data Skew
Partition A:
10 GB
Partition B:
12 GB
Partition C:
950 GB
One partition stores substantially more data than the others.
Traffic Skew
Partition A:
1,000 stored records
100 requests per second
Partition B:
1,000 stored records
100 requests per second
Partition C:
1,000 stored records
20,000 requests per second
Storage is evenly distributed, but request traffic is not.
Analysis rule: Measure both data distribution and traffic distribution. Balanced storage does not prove balanced workload.
Why Hot Partitions Are Dangerous
A hot partition can cause:
- High request latency
- Throttled reads or writes
- Timeouts
- Increased retry traffic
- CPU saturation
- Memory pressure
- Disk input and output saturation
- Network saturation
- Compaction or maintenance backlog
- Queue or stream-consumer lag
- Uneven scaling costs
- Reduced overall throughput
Retries can worsen the hotspot when clients immediately repeat failed requests against the same partition.
Basic Load Model
If one key receives a fraction \(H\) of total requests and total request rate is \(Q\), the request rate for that key is:
\[ HotKeyQPS = Q \times H \]
If the hot key maps to one partition, that partition must absorb the complete hot-key rate in addition to its other requests.
Example
Total request rate:
50,000 requests per second
Share received by one course:
30%
Hot-course traffic:
50,000 × 0.30
=
15,000 requests per second
This is illustrative. Capacity limits depend on record size, operation type, consistency, index maintenance, storage engine, and provider configuration.
Partition Imbalance Ratio
A simple traffic-imbalance indicator is:
\[ ImbalanceRatio = \frac{ BusiestPartitionQPS }{ AveragePartitionQPS } \]
A rising ratio indicates that average cluster utilization is hiding concentrated work.
Evaluate this together with CPU, storage, latency, throttling, and individual operation costs.
Common Causes
Hot partitions commonly result from:
- Low-cardinality partition keys
- Monotonically increasing keys
- Time-only partitioning
- Popular entities or hot keys
- Large tenants
- Global counters
- Popular secondary-index values
- Uneven hash distribution
- Bursty temporal traffic
- Range-oriented routing
- Unbounded partitions
- Incorrect retry behaviour
Cause 1: Low-cardinality Keys
A low-cardinality key has only a small number of possible values.
Poor partition-key candidates:
status
draft
published
archived
difficulty
beginner
intermediate
advanced
isActive
true
false
If most records use published, nearly all relevant writes and
queries can target one logical partition.
Better Direction
Instead of:
status
Consider:
tenantId + status + bucket
Or:
tenantId + courseId
The correct choice depends on
the required query.
Cause 2: Monotonically Increasing Keys
A monotonically increasing key continually produces values at one end of the key range.
Sequential keys:
100001
100002
100003
100004
100005
Timestamp keys:
10:00:01
10:00:02
10:00:03
10:00:04
In a range-partitioned design, new writes can continuously target the newest key range instead of spreading across existing ranges.
Older ranges:
Mostly reads
Newest range:
All current writes
|
v
Hot partition
Cause 3: Time-only Partitioning
Partitioning an event stream only by date or current time can direct all writes for the active period to one partition.
Partition key:
2026-09-23
All events today:
One logical partition
Partition key:
2026-09-23
+
hash(tenantId or entityId)
+
controlled bucket
The second design spreads active-period writes across several partitions while retaining a usable time boundary.
Cause 4: Popular Entities
One course, article, video, product, user, or topic can become much more popular than the average item.
Normal course:
200 reads per minute
Popular course:
100,000 reads per minute
Partition key:
courseId
Result:
The popular course's partition
receives concentrated read traffic.
Hashing the key does not solve a single hot key. Hashing determines which partition receives it, but all requests for that exact key still reach the same partition.
Cause 5: Large Tenants
Partitioning only by tenant can work when tenants are similar in size and traffic. It can fail when one tenant is much larger or busier.
Tenant 1:
10,000 records
Tenant 2:
12,000 records
Tenant 3:
300 million records
Partition key:
tenantId
Result:
Tenant 3 creates a large
and heavily used partition.
Large tenants can require subpartitioning, isolated capacity, or a dedicated placement strategy.
Cause 6: Global Counters
A global counter is one logical record updated by many concurrent writers.
Global value:
totalPageViews
Every page view:
Increment same key
Result:
One record and partition
receive every write.
Global inventory, likes, views, sequence values, and rate-limit counters can become write hotspots.
Cause 7: Hot Secondary-index Keys
The base table can be balanced while a secondary index is hot.
Base table partition key:
courseId
Well distributed:
Many course IDs
Secondary-index partition key:
courseStatus
Most records:
published
Result:
The "published" index partition
receives concentrated activity.
Evaluate base and secondary-index partition keys independently.
Cause 8: Traffic Bursts
A normally balanced system can become temporarily skewed during:
- A new course release
- A live event
- A product launch
- A scheduled batch job
- A notification campaign
- An examination window
- A trending article
- A failover or recovery event
Capacity planning must include peak key-level traffic, not only average table traffic.
Hash vs Range Partitioning
| Area | Hash Partitioning | Range Partitioning |
|---|---|---|
| Placement | Hash of key determines partition | Ordered key range determines partition |
| Distribution | Can distribute varied keys well | Can skew when active keys cluster in one range |
| Range queries | Can require scatter-gather | Naturally supports ordered key ranges |
| Single hot key | Still maps to one partition | Still belongs to one range |
| Sequential writes | Can spread when key identities vary | Can target the newest range repeatedly |
Hashing rule: Hashing helps distribute many different keys. Hashing cannot divide the traffic of one logical hot key unless the application deliberately creates several physical keys.
Symptoms of a Hot Partition
Common symptoms include:
- One partition has much higher request volume
- Only some tenants or keys experience timeouts
- High-percentile latency rises while average latency looks acceptable
- One node has high CPU while others remain underused
- One partition experiences throttling
- Retry traffic concentrates on the same keys
- Compaction or flushing falls behind on one node
- Queue lag appears only for selected partitions
- Adding cluster capacity provides little improvement
- One secondary-index value dominates request metrics
Detecting Hot Partitions
Table-level averages can hide skew. Observe metrics by partition, key range, tenant, and query pattern.
Compare:
- Requests per partition
- Bytes per partition
- CPU per partition host
- Read and write latency per partition
- Throttling per partition
- Storage input and output per partition
- Cache hit rate per partition
- Compaction backlog per partition
- Top keys by request volume
- Top tenants by traffic and stored data
Conceptual SQL Analysis
SELECT
partition_key,
COUNT(*) AS request_count,
SUM(bytes_processed) AS bytes_processed,
AVG(duration_ms) AS average_duration_ms
FROM request_metrics
WHERE recorded_at >= :start_time
AND recorded_at < :end_time
GROUP BY
partition_key
ORDER BY
request_count DESC;
Avoid placing sensitive raw keys in broadly accessible logs. Use approved, protected identifiers or controlled hashes where appropriate.
Top-key Analysis
Identify which logical keys contribute most traffic.
Top key traffic:
course-42 31%
course-81 8%
course-19 4%
all others 57%
This distribution shows that one key, not the overall record count, is the main scaling constraint.
Strategy 1: Choose High-cardinality Keys
A high-cardinality partition key has many possible values and can distribute records more widely.
Partition key:
courseStatus
Values:
draft
published
archived
Partition key:
tenantId + courseId
High cardinality alone is insufficient if one key receives most traffic. Distribution must be evaluated using both stored data and actual access frequency.
Strategy 2: Hash the Partition Key
Hashing can spread different logical keys across partitions.
Logical key:
tenant-17:course-42
Hash:
hash(tenant-17:course-42)
Partition:
hash result maps to one partition
This prevents adjacent or sequential logical keys from automatically landing in the same range.
It does not split one hot logical key across several partitions.
Strategy 3: Time Bucketing
Time bucketing limits the amount of data stored under one time-scoped partition key.
Unbounded partition:
tenant-17:course-42
Daily partitions:
tenant-17:course-42:2026-09-21
tenant-17:course-42:2026-09-22
tenant-17:course-42:2026-09-23
Time buckets can make retention and time-range queries easier, but all active writes can still concentrate in the current bucket.
Time bucketing is often combined with another distribution attribute.
Strategy 4: Write Sharding or Salting
Write sharding creates several physical partition keys for one logical key.
Logical key:
course-42:views
Physical keys:
course-42:views:shard-0
course-42:views:shard-1
course-42:views:shard-2
course-42:views:shard-3
Writes are distributed across the shards using a random or deterministic routing rule.
A total read must query and combine all shards.
| Benefit | Cost |
|---|---|
| Spreads write traffic | Reads can require scatter-gather |
| Reduces single-key contention | Aggregation becomes more complex |
| Supports parallel updates | Strongly consistent totals become harder |
| Can target only exceptional hot keys | Routing metadata or logic is required |
Deterministic Shard Selection
If \(S\) is the number of write shards, one possible routing formula is:
\[ ShardNumber = hash(DistributionValue) \bmod S \]
A stable distribution value can be a request ID, learner ID, device ID, or another field appropriate to the access pattern.
PHP Write-shard Example
<?php
declare(strict_types=1);
function selectWriteShard(
string $distributionValue,
int $shardCount
): int {
if ($shardCount < 1) {
throw new InvalidArgumentException(
'The shard count must be positive.'
);
}
$hash =
hash(
'sha256',
$distributionValue
);
$prefix =
substr(
$hash,
0,
8
);
$number =
hexdec(
$prefix
);
return $number % $shardCount;
}
The shard count, hash function, migration process, and routing contract must be versioned. Changing them without a migration plan can make existing data difficult to locate.
Strategy 5: Composite Partition Keys
A composite key combines fields to achieve useful locality and better distribution.
Access pattern:
Retrieve activity for one course
during one day.
Candidate partition key:
tenantId
+
courseId
+
activityDate
+
writeBucket
The time field bounds partition growth. The write bucket spreads active writes.
Reading the complete day requires querying all configured buckets.
Strategy 6: Cache Hot Reads
Frequently requested, safely cacheable data can be served from a distributed cache, edge cache, or application cache.
Request for popular course
|
v
Check cache
|
+-- Hit:
| return cached course
|
+-- Miss:
read database
populate cache
return result
Caching reduces repeated database reads but requires:
- Expiration policy
- Invalidation strategy
- Tenant-aware keys
- Stampede protection
- Authorization-safe content
- Staleness limits
Strategy 7: Request Coalescing
Request coalescing allows concurrent cache misses for one key to share one backend load.
1,000 simultaneous requests
for course 42
|
v
Cache miss
|
v
One request loads database record
|
v
Other requests wait for same result
|
v
Result cached and shared
Without coalescing, a cache expiration can send a burst of identical requests to the already hot database partition.
Strategy 8: Staggered Expiration
If many cached records expire at the same exact time, their reloads can create a synchronized traffic spike.
Base cache lifetime:
10 minutes
Jittered lifetime:
10 minutes
plus a small randomized variation
Expiration jitter spreads refresh work over time. The permitted variation must remain within the freshness requirement.
Strategy 9: Preaggregation
High-volume events can be aggregated before updating a global total.
Individual view events
|
v
Per-worker or per-shard counters
|
v
Periodic aggregation
|
v
Course view summary
This reduces write contention at the cost of delayed totals.
Counter-shard Model
CREATE TABLE course_view_counter_shards
(
tenant_id BIGINT NOT NULL,
course_id BIGINT NOT NULL,
counter_date DATE NOT NULL,
shard_number INT NOT NULL,
view_count BIGINT NOT NULL,
PRIMARY KEY
(
tenant_id,
course_id,
counter_date,
shard_number
)
);
The query sums all shard values for the requested course and date.
Strategy 10: Isolate Large Tenants
A very large tenant can be assigned dedicated or subdivided capacity.
Normal tenants:
Shared partitioning pool
Large tenant:
Dedicated shard group
or tenant-specific subshards
Routing metadata identifies where each tenant's data belongs.
This approach introduces placement, migration, recovery, monitoring, and operational complexity.
Strategy 11: Split Hot Partitions
Some distributed databases can split an overloaded or oversized partition into smaller ranges.
Before:
Partition A
Keys 1 through 1,000,000
After split:
Partition A1
Keys 1 through 500,000
Partition A2
Keys 500,001 through 1,000,000
Splitting helps when traffic is distributed across many keys in the original partition.
Splitting does not fully solve one indivisible hot key because that key still belongs to one resulting partition.
Strategy 12: Repartitioning
Repartitioning moves existing data to a new key or placement scheme.
Old key:
tenantId
New key:
tenantId + entityBucket
Migration:
1. Introduce new routing version.
2. Backfill existing records.
3. Support reads during transition.
4. Route new writes safely.
5. Verify data completeness.
6. Switch authoritative reads.
7. Retire old placement.
This is a suggested migration pattern. The exact procedure depends on the selected database, consistency requirements, and migration capabilities.
Strategy 13: Rate Limiting
Per-key, per-tenant, or per-partition rate limiting can prevent one workload from exhausting shared capacity.
Incoming request
|
v
Determine trusted tenant and resource key
|
v
Check permitted request rate
|
+-- Allowed:
| process request
|
+-- Exceeded:
reject, delay,
or queue according to policy
Rate limiting protects the system but does not remove the underlying skew.
Strategy 14: Load Shedding
During overload, nonessential work can be rejected, delayed, or served from a degraded path.
Possible degraded responses include:
- Stale cached data within an approved bound
- Reduced result detail
- Delayed analytics update
- Queued background processing
- Explicit rate-limit response
Do not silently degrade operations requiring strong correctness, such as a payment, inventory reservation, or authorization change.
Strategy 15: Backpressure
Backpressure slows producers when downstream partitions cannot process work fast enough.
Producer
|
v
Partitioned queue
|
v
Consumer for hot partition
|
v
Processing capacity reached
|
v
Producer slows,
buffers within limits,
or receives rejection
Without backpressure, queues, memory, retries, and pending work can increase without control.
Retry Storms
Poor retry behaviour can amplify a hotspot.
Hot partition slows
|
v
Requests time out
|
v
Clients retry immediately
|
v
Request volume increases
|
v
Partition slows further
Safer retries use:
- Exponential backoff
- Randomized jitter
- Retry limits
- Operation deadlines
- Idempotency
- Server retry guidance where supported
- Circuit breaking or admission control where appropriate
Mitigation Comparison
| Mitigation | Best For | Main Trade-off |
|---|---|---|
| Hash partitioning | Distributing many different keys | Ordered range queries can become distributed |
| Time bucketing | Bounding event partitions | The current bucket can remain hot |
| Write sharding | One logical write-heavy key | Reads must combine several shards |
| Caching | Repeated reads of popular data | Invalidation and staleness |
| Preaggregation | High-volume counters and metrics | Totals can be delayed |
| Dedicated placement | Very large tenants | Routing and operational complexity |
| Partition splitting | Many hot keys in one range | Does not split one indivisible hot key |
| Rate limiting | Protecting shared capacity | Some requests are delayed or rejected |
Designing a Better Key
A good partition key balances several competing goals.
| Goal | Question |
|---|---|
| Cardinality | Are there enough distinct key values? |
| Distribution | Are data and requests reasonably spread? |
| Locality | Can important queries retrieve related records together? |
| Bounded growth | Can one partition grow without limit? |
| Stable identity | Does the key avoid mutable display values? |
| Peak resilience | Can one popular entity overload its partition? |
| Migration | Can routing evolve when traffic changes? |
Learning-platform Example
Access Pattern
Store lesson activity events.
Retrieve events by:
- Tenant
- Course
- Day
- Time range
Weak Partition Key
Partition key:
activityDate
Problem:
Every course and learner
writes to today's partition.
Improved Candidate
Partition key:
tenantId
+
courseId
+
activityDate
+
writeBucket
Sort key:
activityTime
+
activityId
The tenant, course, and date preserve useful query locality. The write bucket can distribute activity for unusually busy courses.
Reading all activity for a course and date requires querying the configured buckets and merging results by time.
Conceptual Activity Schema
CREATE TABLE course_activity
(
tenant_id BIGINT,
course_id BIGINT,
activity_date DATE,
write_bucket INT,
activity_time TIMESTAMP,
activity_id VARCHAR(100),
learner_id BIGINT,
activity_type VARCHAR(50),
payload TEXT,
PRIMARY KEY
(
(
tenant_id,
course_id,
activity_date,
write_bucket
),
activity_time,
activity_id
)
);
This is conceptual wide-column-style syntax. Exact key rules and query behaviour depend on the selected database.
Choosing the Bucket Count
More buckets distribute writes more widely but make complete reads more expensive.
A bucket-count decision should consider:
- Peak writes for one logical key
- Per-partition throughput
- Expected record size
- Read frequency
- Range-query width
- Aggregation cost
- Concurrency
- Operational limits
Avoid selecting an arbitrary large bucket count. Begin with measured demand and maintain a versioned expansion strategy.
Versioned Bucket Configuration
{
"logicalResource": "tenant-17:course-42:activity",
"routingVersion": 2,
"bucketCount": 8,
"effectiveFrom": "routing-change-time"
}
Routing metadata allows exceptional hot entities to use more buckets without forcing the same design on every normal entity.
Reading from Sharded Keys
Read request:
All course-42 activity
for one day
Application:
Query bucket 0
Query bucket 1
Query bucket 2
Query bucket 3
|
v
Merge ordered iterators
|
v
Apply limit and cursor
|
v
Return results
Bound fan-out and maintain deterministic ordering using a stable tiebreaker such as the activity ID.
Pagination across Buckets
Cursor-based pagination across sharded ranges can record progress for each bucket or use a database-supported distributed iterator.
{
"routingVersion": 2,
"bucketPositions": {
"0": "opaque-position-0",
"1": "opaque-position-1",
"2": "opaque-position-2",
"3": "opaque-position-3"
},
"lastActivityTime": "last-returned-time",
"lastActivityId": "last-returned-id"
}
The cursor should remain opaque to clients and include only the information needed by the server's pagination contract.
Dynamic Hot-key Isolation
Normal course:
One logical partition
Traffic increases
beyond approved level
|
v
Mark course as hot
|
v
Assign several write buckets
|
v
Route new writes using
versioned bucket configuration
|
v
Read old and new routing versions
during migration
This adaptive approach avoids increasing read fan-out for every ordinary key.
Tenant Isolation
Partitioning and salting must preserve trusted tenant scope.
Physical key:
TENANT#17
COURSE#42
DATE#2026-09-23
BUCKET#3
Tenant ID must come from authenticated server-side context. A client must not gain access to another tenant merely by changing a partition-key value.
Correctness-sensitive Hot Keys
Some hot keys represent shared mutable state such as inventory, quota, or a financial balance.
Sharding these values is difficult because the system may need a strongly correct total before approving an operation.
Inventory available:
1 unit
Concurrent purchase requests:
Request A
Request B
Request C
Requirement:
At most one request
can reserve the unit.
Caching or approximate counters cannot replace the authoritative conditional update required by this rule.
Possible designs include serialized ownership, conditional writes, reservations, partitioned inventory units, admission control, or queue-based processing. Selection depends on the required correctness and latency.
Observability
Useful metrics include:
- Requests per partition
- Read and write units per partition
- Bytes stored per partition
- Top keys by request rate
- Top tenants by storage and traffic
- Latency by partition and operation
- Throttled requests by partition
- Retries by partition key
- CPU and disk usage by node
- Cache hit rate for hot keys
- Compaction or flush backlog
- Queue lag by partition
- Routing-version distribution
- Scatter-gather fan-out
- Counter-shard imbalance
Alert Conditions
Alert when:
- One partition receives a disproportionate share of requests
- One key exceeds its approved request rate
- Per-partition latency rises above its objective
- Throttling appears only on selected partitions
- One node is saturated while cluster averages remain low
- A current time bucket grows too quickly
- A secondary-index key becomes concentrated
- Retry traffic amplifies an existing hotspot
- Scatter-gather fan-out exceeds its bound
- Repartitioning or migration falls behind
Troubleshooting Workflow
- Identify whether the issue affects reads, writes, or both.
- Compare table-level and per-partition metrics.
- Identify the busiest logical keys.
- Check data-size distribution.
- Check request distribution.
- Inspect the partition-key cardinality.
- Look for monotonically increasing values.
- Inspect secondary-index partition keys.
- Check current time buckets and global counters.
- Check cache misses and synchronized expiration.
- Check retries, backoff, and client timeouts.
- Determine whether the hotspot contains one key or many keys.
- Select caching, splitting, sharding, isolation, or repartitioning accordingly.
- Load-test the candidate mitigation.
- Monitor read fan-out and consistency after the change.
Common Hot-partition Mistakes
Choosing a Key Only for Uniqueness
A unique key can still produce poor traffic distribution or unsuitable query locality.
Using a Low-cardinality Partition Key
Status, Boolean values, and small categories concentrate many records under a few keys.
Using a Timestamp as the Only Key
Current writes can continuously target the newest partition or range.
Assuming Hashing Solves a Single Hot Key
Hashing selects a partition for the key but does not divide one key's request stream.
Monitoring Only Cluster Averages
Healthy average CPU and latency can hide one saturated partition.
Sharding Every Key Prematurely
Unnecessary write buckets increase read fan-out and aggregation complexity for ordinary keys.
Using Too Many Write Buckets
Writes distribute well, but every complete read must query and merge many partitions.
Ignoring Secondary-index Hotspots
The base table can be balanced while an alternative index is heavily skewed.
Retrying Immediately without Jitter
Synchronized retries amplify traffic against the already overloaded partition.
Using Cache without Stampede Protection
Expiration of a popular key can cause many simultaneous database reads.
Adding Nodes without Changing the Key Pattern
One logical key or range can continue to limit performance despite unused capacity elsewhere.
Ignoring Correctness during Counter Sharding
An approximate distributed total can be unsuitable for inventory, financial, quota, or authorization decisions.
Recommended Test Cases
| Test | Expected Evidence |
|---|---|
| Uniform-key traffic | Records and requests distribute within the expected bounds |
| One hot read key | Cache and request-coalescing behaviour are measured |
| One hot write key | Write-sharding or serialization behaviour is measured |
| Monotonic writes | The newest-range hotspot is identified or prevented |
| Low-cardinality key | Skew and throttling behaviour are visible |
| Large tenant | Tenant-specific subpartitioning or placement is validated |
| Current time bucket | Active-period writes remain within tested capacity |
| Secondary-index hotspot | Base and secondary access paths are monitored independently |
| Cache expiration | Jitter and request coalescing prevent a backend surge |
| Retry storm | Backoff, jitter, and retry limits protect the partition |
| Scatter-gather read | Fan-out, merge latency, and pagination remain bounded |
| Repartitioning | Old and new routing schemes return complete data during migration |
| Node failure | Failover does not concentrate traffic beyond remaining capacity |
| Tenant isolation | Physical key routing cannot bypass authorization boundaries |
Hot-partition Best Practices
Recommended Practices
- Design partition keys from both access patterns and workload distribution.
- Prefer keys with sufficient cardinality.
- Avoid low-cardinality values as standalone partition keys.
- Avoid pure sequential or timestamp keys in susceptible range-partitioned write paths.
- Measure data skew and traffic skew independently.
- Monitor per-partition metrics instead of only cluster averages.
- Identify top keys and heavy tenants before production growth.
- Evaluate base-table and secondary-index keys independently.
- Use time buckets to bound partition growth.
- Combine time with a distribution field for active write streams.
- Use write sharding only when one logical key needs additional distribution.
- Keep shard counts measured, bounded, and versioned.
- Cache repeated read-heavy hot keys safely.
- Use request coalescing and expiration jitter.
- Preaggregate high-volume counters when delayed totals are acceptable.
- Isolate exceptional large tenants where justified.
- Use rate limiting, admission control, and backpressure for overload protection.
- Retry with bounded exponential backoff and jitter.
- Load-test realistic skew, not only uniform random traffic.
- Maintain a tested repartitioning and routing-migration strategy.
Practice Exercise
Redesign the activity and popularity counters of your online learning platform to avoid hot partitions.
Requirements
- Generate activity for several tenants and courses.
- Make one course substantially more active than the others.
- Partition first by course ID and observe the hot key.
- Record per-partition traffic and latency.
- Add a daily time bucket.
- Add a controlled write bucket for the hot course.
- Use a deterministic bucket-selection rule.
- Merge bucket results in activity-time order.
- Implement cursor pagination across buckets.
- Shard the course-view counter.
- Aggregate the counter shards.
- Add caching for the popular course page.
- Add request coalescing for cache misses.
- Add retry backoff and jitter.
- Measure read fan-out introduced by the solution.
- Verify tenant isolation.
- Document when an ordinary course becomes eligible for hot-key isolation.
Partition-design Template
| Access Pattern | Candidate Partition Key | Hotspot Risk | Mitigation |
|---|---|---|---|
| Course activity by day | Tenant, course, date, and write bucket | One very popular course | Adaptive controlled write buckets |
| Popular course details | Tenant and course ID | High repeated reads | Cache, coalescing, and jittered expiration |
| Course view counter | Tenant, course, date, and counter shard | Every view updates one key | Sharded counters and periodic aggregation |
| Published courses | Tenant, category, and bounded bucket | Low-cardinality published status | Avoid status as the standalone partition key |
| Tenant activity | Tenant and subpartition key | One exceptionally large tenant | Tenant-specific subsharding or placement |
Frequently Asked Questions
What is a hot partition?
A hot partition is a partition receiving a disproportionate amount of data, read traffic, write traffic, or processing work.
What is a hot key?
A hot key is one logical key receiving substantially more requests than typical keys.
What is partition skew?
Partition skew is an uneven distribution of stored data or workload across partitions.
Does hashing prevent all hot partitions?
No. Hashing can distribute many different keys, but one hot key still maps to one partition unless the logical key is deliberately subdivided.
Why are timestamps risky partition keys?
In a range-oriented design, current writes can continually target the newest timestamp range.
What is write sharding?
Write sharding divides one logical key into several physical keys so writes can be distributed across multiple partitions.
What is salting?
Salting adds a controlled prefix or suffix to create several physical key variants and distribute records or writes.
What is the cost of write sharding?
Complete reads and aggregates can require querying and merging several shard keys.
Can caching solve a hot write key?
Caching primarily reduces repeated reads. Write hotspots normally require sharding, aggregation, serialization, isolation, or admission control.
Can adding more nodes solve a hot partition?
Not necessarily. If one indivisible logical key receives the traffic, additional nodes do not divide that key's workload automatically.
Can a secondary index become hot?
Yes. Its alternative partition key can have a different and more concentrated distribution than the base table.
How should hot partitions be tested?
Use realistic skewed traffic, popular keys, large tenants, burst periods, retries, cache expiration, and sustained load rather than only uniform random requests.
Key Takeaway
A hot partition occurs when one partition receives a disproportionate share of data or requests, limiting a distributed system to the capacity of that partition. Common causes include low-cardinality keys, monotonically increasing values, current time ranges, popular entities, large tenants, global counters, and concentrated secondary-index keys. Hashing spreads different keys but does not divide one indivisible hot key. Use high-cardinality composite keys, bounded time buckets, controlled write sharding, caching, request coalescing, preaggregation, tenant isolation, partition splitting, throttling, and backpressure according to the workload. Every mitigation has a cost, especially read fan-out, staleness, routing complexity, or reduced availability. Measure per-partition traffic, top keys, high-percentile latency, retries, throttling, and storage skew, and keep a versioned repartitioning strategy for workloads whose distribution changes over time.