Table of Contents

    hot partitions

    NOSQL & ACCESS-PATTERN DESIGN

    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
    Partition-key Goal
    support required queries → keep related data together → distribute storage → distribute traffic → avoid unbounded partitions

    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:

    1. Low-cardinality partition keys
    2. Monotonically increasing keys
    3. Time-only partitioning
    4. Popular entities or hot keys
    5. Large tenants
    6. Global counters
    7. Popular secondary-index values
    8. Uneven hash distribution
    9. Bursty temporal traffic
    10. Range-oriented routing
    11. Unbounded partitions
    12. 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.

    Concentrated writes
    Partition key:
    
    2026-09-23
    
    
    All events today:
    
    One logical partition
    Time plus distribution
    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.

    Low cardinality
    Partition key:
    
    courseStatus
    
    
    Values:
    
    draft
    published
    archived
    Higher cardinality
    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

    1. Identify whether the issue affects reads, writes, or both.
    2. Compare table-level and per-partition metrics.
    3. Identify the busiest logical keys.
    4. Check data-size distribution.
    5. Check request distribution.
    6. Inspect the partition-key cardinality.
    7. Look for monotonically increasing values.
    8. Inspect secondary-index partition keys.
    9. Check current time buckets and global counters.
    10. Check cache misses and synchronized expiration.
    11. Check retries, backoff, and client timeouts.
    12. Determine whether the hotspot contains one key or many keys.
    13. Select caching, splitting, sharding, isolation, or repartitioning accordingly.
    14. Load-test the candidate mitigation.
    15. Monitor read fan-out and consistency after the change.

    Common Hot-partition Mistakes

    1

    Choosing a Key Only for Uniqueness

    A unique key can still produce poor traffic distribution or unsuitable query locality.

    2

    Using a Low-cardinality Partition Key

    Status, Boolean values, and small categories concentrate many records under a few keys.

    3

    Using a Timestamp as the Only Key

    Current writes can continuously target the newest partition or range.

    4

    Assuming Hashing Solves a Single Hot Key

    Hashing selects a partition for the key but does not divide one key's request stream.

    5

    Monitoring Only Cluster Averages

    Healthy average CPU and latency can hide one saturated partition.

    6

    Sharding Every Key Prematurely

    Unnecessary write buckets increase read fan-out and aggregation complexity for ordinary keys.

    7

    Using Too Many Write Buckets

    Writes distribute well, but every complete read must query and merge many partitions.

    8

    Ignoring Secondary-index Hotspots

    The base table can be balanced while an alternative index is heavily skewed.

    9

    Retrying Immediately without Jitter

    Synchronized retries amplify traffic against the already overloaded partition.

    10

    Using Cache without Stampede Protection

    Expiration of a popular key can cause many simultaneous database reads.

    11

    Adding Nodes without Changing the Key Pattern

    One logical key or range can continue to limit performance despite unused capacity elsewhere.

    12

    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

    1. Generate activity for several tenants and courses.
    2. Make one course substantially more active than the others.
    3. Partition first by course ID and observe the hot key.
    4. Record per-partition traffic and latency.
    5. Add a daily time bucket.
    6. Add a controlled write bucket for the hot course.
    7. Use a deterministic bucket-selection rule.
    8. Merge bucket results in activity-time order.
    9. Implement cursor pagination across buckets.
    10. Shard the course-view counter.
    11. Aggregate the counter shards.
    12. Add caching for the popular course page.
    13. Add request coalescing for cache misses.
    14. Add retry backoff and jitter.
    15. Measure read fan-out introduced by the solution.
    16. Verify tenant isolation.
    17. 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

    1

    What is a hot partition?

    A hot partition is a partition receiving a disproportionate amount of data, read traffic, write traffic, or processing work.

    2

    What is a hot key?

    A hot key is one logical key receiving substantially more requests than typical keys.

    3

    What is partition skew?

    Partition skew is an uneven distribution of stored data or workload across partitions.

    4

    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.

    5

    Why are timestamps risky partition keys?

    In a range-oriented design, current writes can continually target the newest timestamp range.

    6

    What is write sharding?

    Write sharding divides one logical key into several physical keys so writes can be distributed across multiple partitions.

    7

    What is salting?

    Salting adds a controlled prefix or suffix to create several physical key variants and distribute records or writes.

    8

    What is the cost of write sharding?

    Complete reads and aggregates can require querying and merging several shard keys.

    9

    Can caching solve a hot write key?

    Caching primarily reduces repeated reads. Write hotspots normally require sharding, aggregation, serialization, isolation, or admission control.

    10

    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.

    11

    Can a secondary index become hot?

    Yes. Its alternative partition key can have a different and more concentrated distribution than the base table.

    12

    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.