hotspots
Hotspots
Learn how uneven traffic, storage, processing, and contention create hot keys, hot partitions, hot shards, and hot nodes. Understand why balanced data does not guarantee balanced workload, how range and hash sharding influence hotspots, and how salting, key splitting, caching, replication, request coalescing, isolation, throttling, adaptive routing, and resharding can protect a distributed system.
Introduction
Sharding distributes data and workload across several database nodes. Ideally, each shard receives a manageable share of storage, reads, writes, CPU usage, memory usage, and network traffic.
Distribution is rarely perfect. One key, tenant, partition, shard, node, or geographic region can receive significantly more load than the others. This overloaded area is called a hotspot.
Shard A:
18% of requests
Shard B:
20% of requests
Shard C:
47% of requests
Shard D:
15% of requests
Result:
Shard C is the workload hotspot.
A hotspot limits the capacity of the complete system. Adding more shards does not help when incoming work continues targeting the same overloaded ownership unit.
Core idea: A distributed system scales only when both data and workload can be distributed. Even row counts are not sufficient when one key, tenant, range, or shard receives most of the traffic.
Prerequisites
| # | Prerequisite | Why It Is Needed |
|---|---|---|
| 1 | Range, hash, directory, and geo sharding | Each placement strategy creates different hotspot risks. |
| 2 | Consistent hashing | Consistent hashing distributes key ownership but does not divide one hot key automatically. |
| 3 | Replication | Read replicas can distribute selected reads from a popular partition. |
| 4 | Caching | Caching can prevent repeated reads from reaching an overloaded datastore. |
| 5 | Resharding | Hot ownership ranges or tenants can require splitting or movement. |
| 6 | Rate limiting and overload control | Traffic must remain bounded while structural corrections are applied. |
| 7 | Observability | Hotspots must be identified by key, shard, tenant, route, and resource. |
What Is a Hotspot?
A hotspot is a disproportionately busy or constrained part of a distributed system that limits overall performance or availability.
A hotspot can consume excessive:
- CPU time
- Memory
- Storage capacity
- Disk input and output operations
- Network bandwidth
- Database connections
- Locks
- Queue capacity
- Cache capacity
- Request-processing concurrency
Types of Hotspots
| Hotspot Type | Meaning | Example |
|---|---|---|
| Hot key | One key receives excessive traffic | One popular course is requested repeatedly |
| Hot tenant | One tenant generates substantially more data or requests | One enterprise tenant runs many reports |
| Hot partition | One logical partition receives concentrated activity | All current events are written to today's partition |
| Hot shard | One shard receives more load than peer shards | Several busy tenants are placed on one shard |
| Hot node | One physical server hosts too many busy ownership units | Several popular hash ranges share one database node |
| Hot range | A contiguous key range receives most new operations | Sequential identifiers send all new writes to the final range |
| Hot region | One geographic shard group receives disproportionate demand | One regional group serves most active users |
| Hot index | Many writes repeatedly update one index page or value | A low-cardinality status field receives concentrated updates |
Data Skew vs Workload Skew
| Condition | Meaning |
|---|---|
| Data skew | Some shards store substantially more records or bytes than others. |
| Workload skew | Some shards receive substantially more operations than others. |
| Resource skew | Some nodes consume more CPU, memory, I/O, or connections. |
| Temporal skew | The hotspot changes according to time, events, or workload phase. |
Balanced storage:
Shard A = 25%
Shard B = 25%
Shard C = 25%
Shard D = 25%
Unbalanced traffic:
Shard A = 10%
Shard B = 12%
Shard C = 68%
Shard D = 10%
Measurement rule: Measure data size and traffic independently. Equal record counts do not prove equal processing cost.
Measuring Distribution
Let \(L_i\) represent the measured load on shard \(i\), and let \(n\) represent the total number of shards.
The average load is:
\[ AverageLoad = \frac{ \sum_{i=1}^{n} L_i }{ n } \]
A simple hotspot ratio is:
\[ HotspotRatio = \frac{ MaximumShardLoad }{ AverageShardLoad } \]
This ratio is only one signal. Compare CPU, storage, I/O, connections, latency, and traffic separately because different resources can have different hotspot owners.
Range-sharding Hotspots
Range sharding places contiguous key intervals together.
Shard A:
IDs 1 to 1,000,000
Shard B:
IDs 1,000,001 to 2,000,000
Shard C:
IDs 2,000,001 to 3,000,000
When the key increases sequentially, new records target the final active range.
New ID:
2,900,001
-> Shard C
New ID:
2,900,002
-> Shard C
New ID:
2,900,003
-> Shard C
Shards A and B can remain mostly idle while Shard C receives all current writes.
Common Monotonic Keys
- Auto-incrementing identifiers
- Creation timestamps
- Ordered event sequence numbers
- Time-based keys
- Increasing invoice or transaction numbers
Time-series Hotspot
Historical partitions:
January -> inactive
February -> inactive
March -> inactive
Current partition:
April -> receives every new event
Time-based partitioning supports time-range scans and retention, but one current partition can receive all inserts.
Possible design directions include:
- Hash by entity or device before applying a time range
- Create several current write buckets
- Split the active range earlier
- Roll partitions at an appropriate size or load boundary
- Buffer and batch writes where correctness allows
Hash-sharding Hotspots
Hash sharding generally distributes many distinct keys more evenly than range sharding.
It does not guarantee even traffic when key popularity differs.
Key distribution:
1,000,000 keys spread evenly
Traffic:
course:42
-> 500,000 requests per minute
Other keys:
-> low traffic
Result:
The shard owning course:42
becomes hot.
Hashing the same key repeatedly always produces the same placement. The hash function cannot divide one indivisible key among several owners.
Consistent Hashing and Hotspots
Consistent hashing minimizes ownership movement during node changes. Virtual nodes can improve distribution of hash ranges among physical nodes.
Consistent hashing does not automatically solve:
- One extremely popular key
- One very large tenant
- Several hot virtual ranges placed on one node
- Unequal physical-server capacity
- Uneven operation cost among keys
Weighted virtual nodes can account for different server capacities, but weights must be validated against actual workload characteristics.
Directory-sharding Hotspots
Directory sharding maps tenants or keys explicitly to shards.
tenant-101 -> shard-a
tenant-102 -> shard-a
tenant-103 -> shard-b
tenant-104 -> shard-c
It allows a noisy or large tenant to be moved independently.
Before:
tenant-101 -> shared shard-a
After:
tenant-101 -> dedicated shard-d
Directory sharding offers flexible hotspot correction, but the move still requires controlled copying, synchronization, validation, and routing cutover.
Geo-sharding Hotspots
Geographic populations and traffic are rarely distributed equally.
Region A:
15% of users
Region B:
60% of users
Region C:
25% of users
Geo sharding can improve regional locality while creating a dominant regional shard group.
Possible controls include:
- Adding capacity inside the hot region
- Sub-sharding the region by tenant or user
- Separating reads from writes
- Using regional cache and delivery layers
- Moving only when approved residency rules permit it
Hot Tenant
Tenant-based sharding colocates tenant data and transactions, but tenants can have dramatically different sizes.
Tenant A:
10,000 records
Tenant B:
25,000 records
Tenant C:
500,000,000 records
Hashing each tenant ID once does not divide Tenant C. The complete tenant remains assigned to one shard.
Large-tenant Directions
- Move the tenant to a dedicated shard
- Sub-shard the tenant by user, account, object, or another key
- Separate interactive and analytical workloads
- Archive cold tenant data
- Apply tenant-specific rate and concurrency policies
Low-cardinality Shard Key
A shard key with few distinct values cannot distribute data among many shards effectively.
Shard key:
status
Possible values:
active
inactive
pending
Available shards:
20
Three key values cannot create balanced ownership across twenty independent shard destinations without another distribution dimension.
Write-amplification Hotspot
One logical write can trigger several physical operations.
One enrollment update
|
+-- Base table update
+-- Secondary-index update
+-- Audit record
+-- Outbox event
+-- Cache invalidation
+-- Search update
A shard with moderate request count can still be hot because its operations are more expensive.
Contention Hotspots
A hotspot can occur around one shared record even when data and traffic are otherwise distributed.
Many workers update:
global_counter
Result:
- Lock contention
- Transaction retries
- Serialized execution
- Increased latency
Examples include:
- Global counters
- One inventory record
- One account balance
- One queue head
- One sequence generator
- One popular metadata record
Detecting Hotspots
Compare metrics at several ownership levels:
- Key
- Tenant
- Partition
- Shard
- Physical node
- Region
- API route
- Database query fingerprint
Architecture hotspot analysis can also correlate recurring incidents, outages, performance problems, bottlenecks, unstable integrations, and overloaded services with affected architecture components.
Conceptual Shard-load Query
SELECT
shard_id,
COUNT(*) AS operation_count,
SUM(duration_ms) AS total_duration_ms,
AVG(duration_ms) AS average_duration_ms,
MAX(duration_ms) AS maximum_duration_ms
FROM database_operation_metrics
WHERE recorded_at >= :window_start
GROUP BY shard_id
ORDER BY total_duration_ms DESC;
Record counts alone are insufficient. Duration, CPU, reads, writes, bytes, locks, and failures can reveal different hotspots.
Hotspot Mitigation Categories
| Category | Purpose |
|---|---|
| Reduce repeated work | Use caching, coalescing, batching, or precomputation |
| Distribute reads | Use replicas, caches, or derived projections |
| Distribute writes | Split keys, add buckets, or change partition ownership |
| Isolate noisy workloads | Use dedicated shards, pools, queues, or resource limits |
| Bound demand | Apply rate, concurrency, payload, and queue limits |
| Change ownership | Split, move, or reshard the hot data |
| Optimize the operation | Improve indexes, queries, storage access, and transaction scope |
Caching Hot Reads
Popular read
|
v
Distributed cache
|
+-- Hit:
| return value
|
+-- Miss:
read authoritative store
populate cache
return value
Caching can protect a datastore from repeated reads when the value is safe to cache for the caller, tenant, authorization scope, and required freshness.
Request Coalescing
Request coalescing lets concurrent misses for one key share one backend operation.
10,000 concurrent requests
for course:42
|
v
One active cache fill
|
v
Other requests wait within
a controlled boundary
|
v
Result is shared
Without coalescing, an expired popular cache entry can create a thundering herd against the backing datastore.
Read Replication
Read replicas or replicated cache entries can distribute reads for a hot key.
Hot read key
|
+-- Replica A
+-- Replica B
+-- Replica C
Account for replication lag, invalidation, authorization, routing, and failure-domain placement.
Read replication does not distribute writes requiring one authoritative order.
Key Salting
Key salting adds a controlled distribution component to a logical key.
Original key:
course:42
Salted keys:
course:42:0
course:42:1
course:42:2
course:42:3
Writes can be distributed among buckets. Reads that need the complete logical value can query and combine all relevant buckets.
Suitable Uses
- Distributed counters
- Append-only events
- Aggregations that can merge bucket results
- Work queues that do not require one strict global order
Risks
- Reads can require fan-out
- Global ordering becomes more difficult
- Uniqueness must be enforced separately
- Transactions spanning buckets require coordination
- Bucket count changes require migration planning
Composite Shard Keys
A composite key can preserve a business ownership prefix while distributing activity across buckets.
Composite key:
tenant_id
+
bucket_number
Example:
tenant-17:bucket-0
tenant-17:bucket-1
tenant-17:bucket-2
tenant-17:bucket-3
The bucket can be derived from a stable record attribute, random assignment, or a hash of a secondary identifier.
Hash Prefix before Time Range
Instead of:
timestamp -> one current partition
Use:
hash(device_id) -> bucket
timestamp -> range inside bucket
This spreads current writes across multiple buckets while retaining time-based organization inside each bucket.
Moving a Hot Tenant
Shared Shard A:
Tenant 11
Tenant 12
Tenant 17
Tenant 19
Tenant 17 becomes hot
|
v
Move Tenant 17
to dedicated Shard D
The move requires a controlled resharding process covering historical data, concurrent changes, validation, versioned routing, cutover, and rollback.
Splitting a Hot Range
Before:
Shard A owns 1 to 10,000,000
After:
Shard A owns 1 to 5,000,000
Shard B owns 5,000,001 to 10,000,000
Select the split point from actual data and workload distribution. Equal key ranges can still produce unequal traffic.
Rate and Concurrency Limits
A structural hotspot cannot always be corrected immediately. Overload controls protect the system while a longer-term fix is implemented.
Incoming operation
|
v
Identify trusted tenant,
key, route, and priority
|
v
Check rate and concurrency
|
+-- Capacity available:
| process
|
+-- Capacity unavailable:
reject, defer,
or queue within a bound
Useful boundaries include:
- Requests per tenant
- Concurrent operations per hot key
- Report generations per shard
- Writes per partition
- Queued jobs per ownership range
Workload Isolation
Interactive API pool:
Course access
Enrollment
Progress updates
Batch worker pool:
Reports
Exports
Analytics
Backfills
Separating expensive batch work can prevent reports, migrations, and analytics from consuming resources required by interactive requests.
Query and Index Optimization
A shard can appear hot because every operation performs unnecessary work.
Review:
- Query plans
- Missing or ineffective indexes
- Large scans
- Repeated joins
- Lock duration
- Transaction size
- Unbounded result sets
- Connection-pool behaviour
- Materialization and sorting
Resharding an inefficient query can distribute the inefficiency instead of removing it.
Diagnosis rule: Identify the constrained resource before adding shards. A CPU-intensive query, global lock, or hot key can remain the bottleneck after resharding.
Autoscaling and Hotspots
Autoscaling adds general capacity. It can help when the workload can be distributed among new instances.
Autoscaling does not automatically help when:
- Every request targets one key
- One database shard owns all relevant data
- One global lock serializes operations
- One external dependency is rate limited
- The shard key does not permit finer distribution
Conceptual Hotspot Policy
hotspotControl:
detection:
dimensions:
- shard
- tenant
- partition
- key
- route
metrics:
- request-rate
- write-rate
- latency
- cpu
- storage-io
- connection-usage
- lock-wait
immediateProtection:
rateLimiting: enabled
concurrencyLimiting: enabled
boundedQueues: enabled
requestCoalescing: enabled
readHotspot:
cache: enabled
replicaReads: workload-specific
writeHotspot:
keySplitting: workload-specific
saltedBuckets: workload-specific
batching: workload-specific
isolation:
hotTenantDedicatedShard: supported
batchWorkloadSeparatePool: true
structuralCorrection:
splitRange: supported
moveTenant: supported
reshardingRequired: validated-process
observability:
hotKeyDetection: enabled
shardSkewDetection: enabled
migrationTracking: enabled
This configuration is conceptual. Thresholds and corrective actions must be based on measured workload, consistency requirements, and database-specific capabilities.
Learning-platform Examples
| Hotspot | Cause | Possible Direction |
|---|---|---|
| Popular course | Many learners request the same content | Cache, replicate, and coalesce concurrent misses |
| Examination-start traffic | Many users access one assessment simultaneously | Pre-warm content, reserve capacity, and bound admission |
| Large enterprise tenant | One tenant exceeds its shared-shard allocation | Move to a dedicated shard or sub-shard the tenant |
| Progress-event partition | All current writes target one time range | Hash by learner or course before applying time ranges |
| Report generation | Expensive reports compete with interactive queries | Use a bounded queue and separate worker or analytical storage |
| Global enrollment counter | Every enrollment updates one record | Use distributed partial counters and controlled aggregation where safe |
Security Considerations
Hotspot controls must preserve tenant and authorization boundaries.
- Derive tenant identity from trusted authentication context
- Do not expose raw hot-key identifiers unnecessarily
- Prevent callers from selecting privileged traffic priority
- Apply rate limits using trusted caller or tenant identity
- Maintain authorization on cached and replicated responses
- Audit tenant moves and shard-directory changes
- Protect hotspot metrics containing customer identifiers
Observability
Useful hotspot metrics include:
- Records and bytes by shard
- Reads and writes by shard
- Requests by tenant and key
- CPU, memory, I/O, and network by node
- Connection usage by shard
- Lock waits by key or table
- Latency by shard and route
- Cache hit rate by key class
- Top-key request concentration
- Queue depth by partition
- Rate-limit and concurrency-limit decisions
- Scatter-gather frequency
- Shard skew and storage skew
- Resharding and tenant-move progress
Structured Hotspot Event
{
"hotspotType": "tenant",
"ownershipUnit": "protected-tenant-reference",
"shardId": "shard-c",
"constrainedResource": "database-connections",
"trafficShare": "measured-value",
"action": "concurrency-limited",
"policyVersion": "approved-policy-version"
}
Avoid including credentials, session identifiers, private request payloads, or unnecessary personal data in hotspot logs.
Alert Conditions
Alert when:
- One shard receives disproportionate traffic
- One tenant exceeds its expected capacity share
- One key dominates request volume
- Shard storage distribution becomes significantly uneven
- One node approaches CPU, memory, I/O, or connection limits
- Lock-wait duration increases around one record or index
- Cache misses for a popular key increase suddenly
- A time-based active partition approaches saturation
- Rate or concurrency limits remain continuously active
- Resharding fails to reduce the original hotspot
- A regional shard group approaches maximum capacity
- Useful throughput falls while incoming traffic remains stable
Troubleshooting Workflow
- Identify the user-visible symptom.
- Identify the constrained resource.
- Compare workload by shard, tenant, partition, and key.
- Separate storage skew from traffic skew.
- Check the shard-key distribution and cardinality.
- Check for monotonically increasing range keys.
- Check for a single popular key or tenant.
- Check query plans, indexes, locks, and transaction duration.
- Check cache hit rate and cache-miss concentration.
- Check replicas and read-routing distribution.
- Check batch jobs, reports, migrations, and background workers.
- Apply immediate rate, concurrency, or queue controls.
- Select caching, splitting, isolation, or resharding based on the cause.
- Verify that the corrective action reduces the constrained-resource load.
Common Hotspot Mistakes
Measuring Row Counts Only
Equal record counts can hide substantial differences in traffic and processing cost.
Using a Sequential Range Key
Every new write targets the shard owning the highest current range.
Assuming Hashing Solves Every Hotspot
One popular key or very large tenant remains concentrated on one owner.
Choosing a Low-cardinality Shard Key
Too few distinct values prevent effective distribution across many shards.
Adding Shards without Changing Ownership
New capacity remains idle while the hot data continues targeting its old shard.
Scaling Compute around a Global Lock
More application instances increase contention against the same serialized resource.
Replicating a Write Hotspot
Read replicas distribute reads but do not necessarily divide authoritative writes.
Salting without Planning Reads
Writes become distributed, but every read requires expensive fan-out and aggregation.
Moving a Hot Tenant to Another Shared Hot Shard
The hotspot is relocated instead of isolated or subdivided.
Ignoring Batch and Analytical Workloads
Reports and exports consume capacity required by interactive operations.
Waiting for Complete Saturation
Emergency resharding must compete with an already overloaded production system.
Applying a Fix without Verifying the Resource
The selected action changes data placement while the true bottleneck remains unchanged.
Recommended Test Cases
| Test | Expected Evidence |
|---|---|
| Uniform key load | Traffic and resource use remain acceptably balanced |
| One hot key | Coalescing, caching, or replication prevents backend collapse |
| One hot tenant | The tenant can be throttled, isolated, moved, or subdivided |
| Sequential write key | The active range hotspot is detected |
| Salted writes | Writes distribute while reads still meet their correctness contract |
| Cache expiration | Request coalescing prevents a thundering herd |
| Global counter | Contention is measured and controlled |
| Heavy report | Batch processing does not exhaust interactive capacity |
| Shard split | The hot range's traffic becomes distributed after cutover |
| Tenant move | The dedicated shard reduces pressure on the original shared shard |
| Hot region | Regional sub-sharding or capacity protects local traffic |
| Overload protection | Rate and concurrency limits preserve useful throughput |
Hotspot Best Practices
Recommended Practices
- Measure workload by key, tenant, partition, shard, node, and region.
- Measure data distribution and traffic distribution separately.
- Choose a high-cardinality shard key.
- Avoid unmodified sequential keys for write-heavy range sharding.
- Use composite keys when one distribution dimension is insufficient.
- Detect hot keys even when hash ranges are balanced.
- Cache popular safe-to-cache reads.
- Use request coalescing for concurrent misses.
- Replicate read-heavy values across appropriate failure domains.
- Apply salting only when results can be recombined safely.
- Isolate large or noisy tenants.
- Separate interactive and analytical workloads.
- Optimize queries and indexes before resharding.
- Apply rate and concurrency limits during overload.
- Use bounded queues and backpressure.
- Begin structural correction before capacity is exhausted.
- Split or move ownership through a validated resharding process.
- Protect tenant and authorization context in every route and cache.
- Verify that remediation reduces the actual constrained resource.
- Test hot keys, tenants, ranges, regions, and recovery behaviour.
Practice Exercise
Design hotspot protection for your online learning platform.
Requirements
- Measure traffic and storage by shard.
- Record the most frequently accessed course keys.
- Identify the largest and busiest tenants.
- Create a traffic-skew dashboard.
- Simulate one highly popular course.
- Add caching and request coalescing.
- Simulate one enterprise tenant generating heavy reports.
- Move reports to a bounded worker queue.
- Apply per-tenant concurrency limits.
- Test a sequential progress-event key.
- Replace it with a composite distributed key where appropriate.
- Move one large tenant to a dedicated shard.
- Validate the move and routing cutover.
- Confirm that the original shard's resource pressure decreases.
- Test traffic recovery after the hotspot ends.
Hotspot-design Template
| Hotspot | Immediate Protection | Structural Correction |
|---|---|---|
| Popular read key | Cache, coalesce, and limit backend concurrency | Replicate or precompute the value |
| Hot write key | Bound write concurrency and queue depth | Split the key into safely mergeable buckets |
| Large tenant | Apply tenant-specific limits | Move to a dedicated shard or sub-shard |
| Current time range | Batch or throttle writes where allowed | Hash by entity before applying time ranges |
| Expensive reports | Use bounded worker concurrency | Move reporting to a separate analytical path |
| Hot geographic region | Add approved regional capacity | Sub-shard the region by tenant or user |
Frequently Asked Questions
What is a hotspot?
A hotspot is a key, tenant, partition, shard, node, or region that receives disproportionately high load or resource consumption.
What is a hot key?
A hot key is one key that receives substantially more operations than surrounding keys.
What is a hot shard?
A hot shard is a shard whose traffic or resource use substantially exceeds that of peer shards.
Why do sequential keys create hotspots?
In range sharding, every new sequential value is routed to the shard owning the latest range.
Does hash sharding prevent hotspots?
It can distribute many keys broadly, but it cannot split one highly popular key or large tenant automatically.
Does consistent hashing solve hot keys?
No. Consistent hashing minimizes remapping during membership changes, but one key still has a primary ownership location.
What is key salting?
Key salting adds a controlled bucket component so one logical workload can be distributed among several physical keys.
Why can salting be expensive for reads?
A complete logical result can require reading and combining multiple buckets.
Can replicas fix a write hotspot?
Replicas can distribute selected reads, but authoritative writes can remain concentrated on one ordered write path.
When should a hot tenant receive a dedicated shard?
Consider dedicated placement when the tenant's storage or workload exceeds the safe shared-shard allocation and cannot be controlled adequately by ordinary limits.
Can autoscaling correct a hotspot?
Only when new capacity can receive part of the workload. It does not fix one indivisible key, global lock, or fixed shard owner automatically.
How should a hotspot be fixed?
First identify the constrained resource. Then reduce repeated work, distribute reads or writes, isolate noisy workloads, bound demand, optimize operations, or reshard the affected ownership unit.
Key Takeaway
A hotspot occurs when one key, tenant, partition, shard, node, or region receives a disproportionate share of workload or resource consumption. Range sharding is vulnerable to monotonically increasing write keys, while hash and consistent-hash placement can distribute many keys but cannot divide one popular key or oversized tenant automatically. Detect hotspots by measuring traffic, latency, connections, locks, CPU, storage, and I/O at several ownership levels. Protect the system immediately with caching, request coalescing, rate limits, concurrency limits, bounded queues, and workload isolation. Apply structural corrections through better shard keys, composite distribution, salting, read replication, dedicated tenant placement, range splitting, or resharding. Finally, verify that the chosen correction reduces the actual constrained resource rather than merely moving the hotspot to another shard.