Table of Contents

    range, hash, directory and geo sharding

    DISTRIBUTED DATA & SHARDING

    Range, Hash, Directory, and Geo Sharding

    Learn how database sharding divides one logical dataset across multiple physical shards, how range, hash, directory, and geographic strategies route records, and how shard-key selection, data locality, hot partitions, scatter-gather queries, cross-shard transactions, rebalancing, replication, global indexes, tenant isolation, failure handling, and resharding influence a production design.

    Introduction

    A single database server has finite storage, CPU, memory, connection, and write-processing capacity. Read replicas can distribute reads, but every replica still stores a copy of the complete dataset, and writes commonly continue through one primary write path.

    Sharding addresses a different scaling problem by dividing one logical dataset into smaller subsets and placing those subsets on different database servers.

    Logical dataset
          |
          +-- Shard A: subset of records
          +-- Shard B: subset of records
          +-- Shard C: subset of records
          +-- Shard D: subset of records

    The complete logical dataset is formed from the records distributed across all shards.

    Core idea: Replication creates additional copies of data. Sharding splits different data across different database partitions. A production system commonly combines both, with every shard replicated for availability.

    Every sharded system must answer one routing question:

    Given a record or request,
    which shard owns the data?

    The four common answers are:

    • Place it according to a key range
    • Place it according to a hash of the key
    • Look up the owning shard in a directory
    • Place it according to a geographic or regulatory boundary

    Prerequisites

    # Prerequisite Why It Is Needed
    1 Relational data modeling Relationships and query patterns determine which records should remain together.
    2 Indexes and query plans The shard router and local database must efficiently locate records.
    3 Transactions Transactions spanning several shards require additional coordination.
    4 Replication Each shard normally requires replicas for availability and recovery.
    5 Consistent hashing Hash-based systems need a strategy for adding and removing shard capacity.
    6 Data residency and geography Geo sharding can route data according to region, latency, or residency requirements.
    7 Observability Shard size, traffic, storage, latency, and skew must be monitored independently.

    What Is Sharding?

    Sharding is horizontal data partitioning. Rows or records from one logical dataset are divided across several physical databases or partitions.

    Users table:
    
    Users 1 to 1,000,000
          |
          +-- Shard 1
          +-- Shard 2
          +-- Shard 3
          +-- Shard 4

    The application, database platform, middleware, or routing service uses a shard key to identify the shard responsible for a record.

    Sharded Request Flow
    extract shard key → resolve owning shard → route operation → execute locally → return or combine result

    Sharding vs Replication vs Partitioning

    Concept Primary Purpose Data Placement
    Replication Availability, durability, read scale, and recovery Several nodes hold copies of the same data
    Sharding Distribute storage and processing across multiple data owners Different shards hold different subsets
    Local table partitioning Manage and access sections of a table within a database platform Partitions can remain under one logical database system

    Sharding with Replication

    Shard A
      +-- Leader
      +-- Replica 1
      +-- Replica 2
    
    
    Shard B
      +-- Leader
      +-- Replica 1
      +-- Replica 2
    
    
    Shard C
      +-- Leader
      +-- Replica 1
      +-- Replica 2

    Sharding decides which shard owns the record. Replication decides how many copies of that shard's data exist and how those copies remain synchronized.

    The Shard Key

    The shard key is the value used to route a record or request to its owning shard.

    Common shard-key candidates include:

    • Tenant ID
    • Account ID
    • User ID
    • Order ID
    • Timestamp
    • Country or region code
    • A composite such as tenant ID and record ID

    A strong shard key should support:

    • Balanced storage distribution
    • Balanced request distribution
    • Efficient routing
    • Common query patterns
    • Data locality for related records
    • Manageable rebalancing
    • Required tenant and security boundaries

    Shard-key rule: Choose the shard key from access patterns and workload distribution, not simply from the easiest available column. A key that distributes rows evenly can still distribute traffic unevenly.

    Hot Shards

    A hot shard receives substantially more traffic or data growth than the other shards.

    Shard A:
    
    10% of requests
    
    
    Shard B:
    
    12% of requests
    
    
    Shard C:
    
    68% of requests
    
    
    Shard D:
    
    10% of requests

    A hot shard can result from:

    • A monotonically increasing key
    • One unusually busy tenant
    • A popular record or partition
    • Uneven geographic traffic
    • Poorly selected range boundaries
    • A low-cardinality shard key
    • Seasonal or event-driven demand

    Range Sharding

    Range sharding divides an ordered keyspace into contiguous intervals. Each interval is assigned to a shard.

    Shard A:
    
    Customer IDs 1 to 999,999
    
    
    Shard B:
    
    Customer IDs 1,000,000 to 1,999,999
    
    
    Shard C:
    
    Customer IDs 2,000,000 to 2,999,999

    Range boundaries can be based on numeric IDs, timestamps, alphabetical keys, dates, prices, or another sortable value.

    Range-routing Flow

    Shard key = 1,450,210
          |
          v
    Compare with configured boundaries
          |
          v
    Matched range:
    
    1,000,000 to 1,999,999
          |
          v
    Route to Shard B

    Range-sharding Advantages

    • Simple conceptual routing
    • Ordered records remain close together
    • Range scans can target one or a few shards
    • Time-based archival can align with shard boundaries
    • Adjacent ranges can be split when they grow

    Range-sharding Limitations

    • Sequential writes can concentrate on the latest range
    • Uneven value distribution creates uneven shard sizes
    • Boundary changes can require data movement
    • A popular range can become a traffic hotspot
    • An initially balanced keyspace can become unbalanced over time

    Time-range Sharding

    Shard 1:
    
    January to March
    
    
    Shard 2:
    
    April to June
    
    
    Shard 3:
    
    July to September
    
    
    Shard 4:
    
    October to December

    Time-range sharding is useful when queries commonly request bounded time intervals and older data becomes less active.

    However, new writes typically target the current time range, potentially creating one hot active shard while historical shards remain mostly idle.

    Conceptual Range Router

    <?php
    
    declare(strict_types=1);
    
    final class RangeShardRouter
    {
        public function findShard(
            int $customerId
        ): string {
            return match (true) {
                $customerId < 1_000_000
                    => 'shard-a',
    
                $customerId < 2_000_000
                    => 'shard-b',
    
                $customerId < 3_000_000
                    => 'shard-c',
    
                default
                    => 'shard-d'
            };
        }
    }

    Production boundaries should be configuration-driven, versioned, validated, and updated through a controlled resharding process rather than hard-coded into every application instance.

    Hash Sharding

    Hash sharding applies a deterministic hash function to the shard key and maps the resulting value to a shard or hash range.

    A simplified direct-modulo model is:

    \[ ShardNumber = Hash(ShardKey) \bmod NumberOfShards \]

    Tenant ID:
    
    tenant-1042
    
    
    Hash result:
    
    8,315,729
    
    
    Shard count:
    
    4
    
    
    Shard:
    
    8,315,729 mod 4
    =
    Shard 1

    Hash-routing Flow

    Extract shard key
          |
          v
    Apply deterministic hash
          |
          v
    Map hash value to shard
          |
          v
    Route request

    Hash-sharding Advantages

    • A suitable hash can distribute keys evenly
    • Sequential IDs are spread across shards
    • Point lookups are efficient when the shard key is available
    • Natural key ordering is less likely to create one current-range hotspot

    Hash-sharding Limitations

    • Adjacent keys can be placed on unrelated shards
    • Range queries can require scatter-gather execution
    • Changing the shard count can move many records with simple modulo hashing
    • A single hot key remains hot even when other keys are balanced
    • Operations without the shard key can require a global lookup

    Conceptual Hash Router

    <?php
    
    declare(strict_types=1);
    
    final class HashShardRouter
    {
        public function __construct(
            private int $shardCount
        ) {
            if ($this->shardCount < 1) {
                throw new InvalidArgumentException(
                    'Shard count must be positive.'
                );
            }
        }
    
        public function findShard(
            string $tenantId
        ): int {
            $unsignedHash =
                (int)sprintf(
                    '%u',
                    crc32($tenantId)
                );
    
            return $unsignedHash
                % $this->shardCount;
        }
    }

    This is a teaching example. Production systems must use the selected platform's supported hashing and placement algorithm. Changing hash algorithms or shard counts without a migration plan can route existing keys incorrectly.

    Modulo Hashing and Resharding

    With direct modulo hashing, changing the shard count changes the result for many keys.

    Before:
    
    hash(key) mod 4
    
    
    After adding a shard:
    
    hash(key) mod 5
    
    
    Result:
    
    Many existing keys map
    to different shard numbers.

    Moving from four to five shards can therefore require substantial data movement.

    Consistent Hashing

    Consistent hashing maps hashes onto a logical ring. Shards own sections of that ring, often through several virtual positions.

    Hash ring:
    
    0 --------------------------- Maximum
         |       |       |       |
       Shard A Shard B Shard C Shard D
    
    
    Key placement:
    
    Hash the key
          |
          v
    Locate its ring position
          |
          v
    Route to the responsible shard

    When a shard is added or removed, the design aims to move only the affected hash ranges rather than recalculating placement over the full keyspace.

    Consistent hashing reduces remapping but does not remove the operational need to copy data, track ownership, preserve availability, and complete migration.

    Directory Sharding

    Directory sharding stores an explicit mapping between a shard key and the shard that owns it.

    Shard directory:
    
    tenant-101 -> shard-a
    tenant-102 -> shard-c
    tenant-103 -> shard-b
    tenant-104 -> shard-a

    Any tenant or key can be assigned to any compatible shard. The routing decision comes from the directory rather than a fixed range or hash formula.

    Directory-routing Flow

    Request contains tenant ID
          |
          v
    Query or read cached directory entry
          |
          v
    Resolve owning shard
          |
          v
    Route operation to shard

    Directory-sharding Advantages

    • Flexible key-to-shard placement
    • One hot or large tenant can be moved independently
    • Shard assignment does not need to follow a mathematical formula
    • Different tenants can receive different placement policies
    • Incremental rebalancing can be performed key by key or group by group

    Directory-sharding Limitations

    • The directory becomes critical routing infrastructure
    • A directory lookup or cache is required before routing
    • Stale directory entries can send requests to the wrong shard
    • Directory updates and data movement must be coordinated
    • Directory availability and recovery require explicit design

    Conceptual Shard Directory

    CREATE TABLE shard_directory
    (
        tenant_id BIGINT PRIMARY KEY,
        shard_id VARCHAR(50) NOT NULL,
        placement_version BIGINT NOT NULL,
        placement_status VARCHAR(30) NOT NULL,
        updated_at TIMESTAMP NOT NULL
    );

    A placement version helps routers recognize stale cached mappings. The migration status can indicate whether a tenant is stable, copying, validating, or being switched.

    Directory Lookup Example

    <?php
    
    declare(strict_types=1);
    
    final class DirectoryShardRouter
    {
        public function __construct(
            private ShardDirectory $directory
        ) {
        }
    
        public function findShard(
            int $tenantId
        ): ShardPlacement {
            $placement =
                $this->directory->findByTenantId(
                    $tenantId
                );
    
            if ($placement === null) {
                throw new RuntimeException(
                    'No shard placement exists for the tenant.'
                );
            }
    
            if (!$placement->isActive()) {
                throw new RuntimeException(
                    'The shard placement is not active.'
                );
            }
    
            return $placement;
        }
    }

    Geo Sharding

    Geo sharding assigns data to shards according to a geographic, jurisdiction, residency, or regional ownership key.

    Region key:
    
    IN
        -> India shard group
    
    
    EU
        -> European shard group
    
    
    US
        -> United States shard group
    
    
    APAC
        -> Asia-Pacific shard group

    The routing key can be based on a trusted account region, tenant residency policy, business location, or another approved ownership attribute.

    Geo-routing Flow

    Request arrives
          |
          v
    Authenticate caller
          |
          v
    Resolve trusted tenant
    or account region
          |
          v
    Select regional shard group
          |
          v
    Route operation

    Geo-sharding Advantages

    • Data can be placed closer to its primary users
    • Regional requests can avoid unnecessary cross-region data access
    • Placement can align with documented data-residency rules
    • Regional failures can be isolated when the architecture supports it
    • Regional teams or services can own local data paths

    Geo-sharding Limitations

    • Regional traffic and storage can be highly uneven
    • Users or tenants can move between regions
    • Global queries require cross-region fan-out or separate global projections
    • Cross-region transactions are expensive and failure-sensitive
    • Residency, backup, replication, and disaster-recovery rules require careful alignment
    • A geographic attribute supplied by the public client cannot be trusted automatically

    Geo-routing rule: Derive regional ownership from trusted server-side account, tenant, contractual, or policy data. Do not let a caller move protected records between jurisdictions by changing a request parameter.

    Geo Sharding vs Geospatial Querying

    Concept Purpose Example
    Geo sharding Decides which physical shard owns the record A European tenant is assigned to a European shard group
    Geospatial query Finds objects using coordinates, distance, points, or polygons Find all learning centers inside a selected region

    A system can use geo sharding for regional ownership and geospatial indexes inside each shard for location-based searches.

    Complete Strategy Comparison

    Area Range Hash Directory Geo
    Placement rule Contiguous key interval Hash-derived shard or token Explicit key-to-shard mapping Trusted regional ownership
    Point lookup Efficient with shard key Efficient with shard key Requires directory resolution Efficient with region and local key
    Range query Efficient when range matches the shard key Often requires scatter-gather Depends on mapped key distribution Efficient inside one region, global query can fan out
    Distribution Depends on range boundaries Normally balanced for many well-distributed keys Controlled explicitly Depends on regional traffic and data
    Hotspot risk High for sequential or popular ranges Lower for many keys, but hot keys remain possible Can isolate hot tenants manually High when one region dominates
    Rebalancing Move or split key ranges Move hash ranges or tokens Update mapping and move selected keys Move regional ownership under strict controls
    Routing dependency Boundary metadata Hash and ownership metadata Highly available directory Trusted region and regional routing metadata
    Typical strength Ordered locality Broad distribution Placement flexibility Regional locality and residency

    Point Queries

    A point query retrieves a specific record using its shard key.

    SELECT
        course_id,
        title,
        status
    FROM courses
    WHERE tenant_id = :tenant_id
      AND course_id = :course_id;

    If tenant_id is the shard key, the router can identify one shard before executing the query.

    Scatter-Gather Queries

    A scatter-gather query sends work to several or all shards and combines their results.

    Global query
          |
          +-- Query Shard A
          +-- Query Shard B
          +-- Query Shard C
          +-- Query Shard D
          |
          v
    Merge, sort, aggregate,
    or paginate results

    Scatter-gather can increase:

    • Query latency
    • Coordinator CPU and memory
    • Network traffic
    • Failure probability
    • Pagination complexity
    • Load on otherwise unrelated shards

    Query rule: The most scalable sharded query contains the shard key and reaches one shard. Treat frequent scatter-gather execution as a signal to revisit the data model, index, projection, or shard key.

    Global Sorting and Pagination

    Pagination becomes harder when records are distributed among shards.

    Shard A returns top results
    
    Shard B returns top results
    
    Shard C returns top results
    
    Coordinator:
    
    - Merges result streams
    - Applies global ordering
    - Selects page boundary
    - Produces continuation state

    Offset pagination can become inefficient and unstable across changing shards. Cursor-based pagination can carry per-shard continuation state or a global ordering key, depending on the data model.

    Global Secondary Indexes

    A query may not contain the shard key.

    Primary placement:
    
    tenant_id -> shard
    
    
    Lookup request:
    
    Find user by email address
    
    
    Problem:
    
    Email alone does not reveal
    the owning tenant shard.

    Possible design directions include:

    • A global lookup index mapping email to tenant and shard
    • A directory service
    • A separate searchable projection
    • A globally unique identifier containing routing information
    • A controlled scatter-gather query for infrequent administration

    A global index introduces its own consistency, availability, uniqueness, repair, and migration requirements.

    Data Colocation

    Records commonly accessed or updated together should be colocated where practical.

    Tenant 17 shard:
    
    - Tenant profile
    - Users
    - Courses
    - Enrollments
    - Progress
    - Tenant configuration

    Sharding related tables by the same tenant key can keep many transactions local to one shard.

    Data belonging to several tenants or global reference data requires a separate strategy.

    Cross-Shard Transactions

    A cross-shard transaction reads or writes data owned by more than one shard.

    Operation:
    
    Move data from Tenant A
    to Tenant B
    
    
    Shard A:
    
    Remove or update source
    
    
    Shard B:
    
    Create or update destination

    Cross-shard transactions can require distributed coordination, increase latency, and expose partial-failure scenarios.

    Possible approaches include:

    • Redesign the ownership boundary to keep the operation local
    • Use database-supported distributed transactions
    • Use an idempotent workflow or saga
    • Use an outbox and durable events
    • Define compensating actions

    Transaction rule: Design the shard key so the most frequent correctness-sensitive transactions stay within one shard. Distributed transactions should be explicit exceptions, not accidental consequences of routine requests.

    Global Uniqueness

    A unique constraint inside one shard does not automatically guarantee uniqueness across every shard.

    Shard A accepts:
    
    username = learner42
    
    
    Shard B also accepts:
    
    username = learner42
    
    
    Each local unique index succeeds.

    Possible directions include:

    • Scope uniqueness by tenant or shard
    • Use globally unique generated identifiers
    • Maintain a strongly coordinated global registry
    • Reserve namespaces or identifier ranges
    • Detect and resolve duplicates according to business policy

    Tenant-based Sharding

    Multi-tenant applications commonly use tenant_id as the shard key.

    Benefits

    • Tenant data can remain colocated
    • Tenant-scoped queries target one shard
    • Tenant migration can be handled as an ownership move
    • Isolation and placement policy can align with the tenant boundary

    Risks

    • A very large tenant can exceed one shard's capacity
    • A busy tenant can create a hot shard
    • Global reporting requires cross-shard aggregation
    • Tenant moves require coordinated directory and data migration

    Large-tenant Strategy

    Small and medium tenants:
    
    Shared tenant shards
    
    
    Large tenant:
    
    Dedicated shard
    
    
    Extremely large tenant:
    
    Sub-sharded by tenant_id
    plus another high-cardinality key

    Directory sharding is useful when selected tenants need dedicated placement without changing every other tenant's assignment.

    Resharding

    Resharding changes ownership of part of the dataset.

    Source shard owns key range
          |
          v
    Create target shard
          |
          v
    Copy historical data
          |
          v
    Capture new changes
          |
          v
    Bring target current
          |
          v
    Validate target
          |
          v
    Switch routing
          |
          v
    Drain old ownership
          |
          v
    Remove old copy after retention

    Exact procedures depend on the selected database. The design must define how writes are handled while data is moving.

    Resharding Write Strategies

    Strategy General Direction Main Risk
    Pause writes Temporarily stop modifications while ownership changes Availability interruption
    Dual write Write to source and target during migration Partial success and divergence
    Change stream Copy base data and replay later changes Lag, ordering, retention, or capture gaps
    Forwarding Old owner forwards operations to the new owner Temporary routing complexity and loops
    Database-managed move Use the platform's supported online resharding mechanism Product limits, capacity pressure, and operational prerequisites

    Shard Split vs Shard Move

    Operation Meaning
    Shard split One shard's keyspace is divided into smaller ownership units
    Shard move An existing ownership unit is transferred to another node or shard group
    Shard merge Several small compatible ownership units are combined
    Shard rebalancing Ownership is redistributed to improve capacity or traffic balance

    Shard Routing and Security

    Routing must not replace authorization.

    Client sends:
    
    tenantId = 17
    
    
    Application must:
    
    1. Authenticate caller
    2. Determine authorized tenant context
    3. Compare route tenant with trusted context
    4. Resolve shard
    5. Authorize requested resource
    6. Execute query

    A caller must not access another tenant's shard merely by changing a tenant identifier in the URL or request body.

    Directory Security

    The shard directory contains sensitive infrastructure routing information.

    Protect it with:

    • Least-privilege access
    • Private network paths
    • Authenticated and authorized updates
    • Versioned placement records
    • Audited ownership changes
    • Encrypted communication and storage according to policy
    • Backup and tested recovery

    Shard Failure

    A shard can become unavailable while other shards continue serving traffic.

    Shard A:
    
    Healthy
    
    
    Shard B:
    
    Unavailable
    
    
    Shard C:
    
    Healthy
    
    
    Impact:
    
    Only data owned by Shard B
    is unavailable when isolation
    and routing are working correctly.

    Each shard should have an appropriate replication and failover strategy. Shard-level blast-radius isolation is useful only when shared routing, directory, identity, and network components remain available.

    Caching Shard Mappings

    Applications can cache range, token, directory, or geo-placement metadata to avoid a central lookup for every request.

    Request
       |
       v
    Check local mapping cache
       |
       +-- Fresh mapping:
       |      route request
       |
       +-- Missing or stale:
              fetch current mapping
              update cache
              route request

    The router needs a way to detect stale mappings after a tenant or key range moves.

    Possible controls include:

    • Placement versions
    • Bounded cache lifetime
    • Invalidation notifications
    • Wrong-owner responses containing safe redirect metadata
    • A lookup retry after routing failure

    Choosing a Sharding Strategy

    Choose Range Sharding When

    • Range scans are a primary access pattern
    • Ordered locality is valuable
    • Boundaries can be split and rebalanced safely
    • Monotonic-key hotspots can be controlled

    Choose Hash Sharding When

    • Point lookups dominate
    • Broad key distribution is important
    • Natural ordering is not needed for most queries
    • Scatter-gather range queries remain uncommon

    Choose Directory Sharding When

    • Placement must be controlled independently for each tenant or key
    • Large tenants need dedicated shards
    • Selective movement is more important than formula-based routing
    • A highly available placement directory can be operated reliably

    Choose Geo Sharding When

    • Regional latency is a primary requirement
    • Data ownership naturally follows a geographic boundary
    • Documented residency constraints influence placement
    • Cross-region queries and transactions are limited or explicitly designed

    Hybrid Sharding

    Production systems frequently combine strategies.

    Step 1:
    
    Geo shard by residency region
    
    
    Step 2:
    
    Directory-map tenant to
    a shard group in that region
    
    
    Step 3:
    
    Hash-distribute tenant records
    inside the assigned shard group
    
    
    Step 4:
    
    Range-partition historical events
    inside each physical shard

    Hybrid designs can solve several placement requirements but increase routing, migration, monitoring, and operational complexity.

    Learning-platform Examples

    Data Possible Strategy Reasoning
    Tenant courses and enrollments Directory or hash sharding by tenant ID Tenant-scoped transactions and queries remain colocated
    Activity-event history Time range combined with a distribution prefix Supports time retention while reducing one current-range hotspot
    Global public course search Separate searchable projection A global search should not repeatedly scan transactional shards
    Large enterprise tenant Directory mapping to a dedicated shard Prevents one tenant from overloading a shared shard
    Region-controlled learner data Geo shard followed by tenant placement Routes data according to approved regional ownership
    Public analytics Asynchronous warehouse or analytical projection Avoids repeated transactional scatter-gather operations

    Conceptual Sharding Policy

    sharding:
      primaryKey:
        type: trusted-tenant-id
    
      placement:
        strategy: directory
    
      directory:
        cacheEnabled: true
        placementVersionRequired: true
        unavailablePolicy: controlled-failure
    
      tenantPolicies:
        default:
          placement: shared-shard-group
    
        largeTenant:
          placement: dedicated-shard
    
        regulatedTenant:
          placement: approved-regional-shard-group
    
      migration:
        copyHistoricalData: true
        captureConcurrentChanges: true
        validateBeforeCutover: true
        rollbackPlanRequired: true
    
      queries:
        crossShard:
          restricted: true
    
        globalAnalytics:
          target: analytical-projection
    
      observability:
        shardSize: enabled
        shardTraffic: enabled
        routingFailures: enabled
        hotKeyDetection: enabled
        migrationProgress: enabled

    This is conceptual configuration. Product-specific routing, resharding, replication, transaction, and failover guarantees must be verified against the selected database.

    Observability

    Useful sharding metrics include:

    • Storage size by shard
    • Read and write rate by shard
    • CPU, memory, and I/O by shard
    • Connection usage by shard
    • Request latency by shard
    • Rows or documents by shard
    • Hot-key and hot-tenant traffic
    • Scatter-gather query count
    • Cross-shard transaction count
    • Shard-directory lookup latency
    • Routing-cache misses
    • Wrong-owner or stale-mapping responses
    • Resharding copy and catch-up progress
    • Shard replication lag
    • Regional traffic and storage distribution

    Structured Routing Event

    {
      "shardKeyType": "tenant",
      "shardId": "regional-shard-group-3",
      "placementVersion": 42,
      "routingStrategy": "directory",
      "operationType": "read",
      "result": "routed"
    }

    Avoid writing confidential tenant data, credentials, database connection details, or raw routing secrets to routine logs.

    Alert Conditions

    Alert when:

    • One shard receives disproportionate traffic
    • Shard storage distribution becomes significantly uneven
    • A shard approaches its capacity boundary
    • Scatter-gather query volume increases unexpectedly
    • Cross-shard transaction failures increase
    • The shard directory becomes unavailable
    • Routing-cache entries repeatedly become stale
    • Requests are routed to the wrong owner
    • Resharding stops or falls behind incoming changes
    • One regional shard group approaches saturation
    • Shard replication falls below its safe requirement
    • A tenant exceeds the safe capacity of its assigned shard

    Troubleshooting Workflow

    1. Identify the sharding strategy and shard key.
    2. Extract the trusted shard-key value from the failing request.
    3. Resolve the expected owning shard.
    4. Compare the router's placement version with the current metadata.
    5. Check the target shard's health and replication state.
    6. Check CPU, memory, storage, connections, and I/O by shard.
    7. Check for hot tenants, keys, ranges, or regions.
    8. Check whether the query contains the shard key.
    9. Check for accidental scatter-gather execution.
    10. Check global index or directory consistency.
    11. Check active data migrations and dual-write behaviour.
    12. Check cross-shard transaction and retry state.
    13. Check recent boundary, token, directory, or regional placement changes.
    14. Rebalance or reshard through the database's approved procedure.

    Common Sharding Mistakes

    1

    Sharding before the Bottleneck Is Proven

    The system accepts substantial routing and transaction complexity when indexing, caching, query tuning, archival, or read replicas could address the actual problem.

    2

    Choosing a Low-cardinality Shard Key

    Too few distinct values prevent meaningful distribution across shards.

    3

    Using a Sequential Key with Fixed Ranges

    New writes concentrate on the shard owning the latest key range.

    4

    Assuming Hashing Eliminates Every Hotspot

    One popular key or tenant remains hot even when other keys are evenly distributed.

    5

    Changing the Shard Count with Simple Modulo Hashing

    Many keys suddenly map to different shards, requiring extensive data movement.

    6

    Making the Directory a Single Failure Point

    Healthy data shards become unreachable because routing metadata is unavailable.

    7

    Trusting a Client-provided Region

    A caller can attempt to route protected records to another jurisdiction or tenant shard.

    8

    Ignoring Cross-shard Query Cost

    Scatter-gather reads consume resources on every shard and create expensive merge operations.

    9

    Assuming Local Uniqueness Is Global

    Different shards can accept the same supposedly unique value.

    10

    Using Unsafe Dual Writes during Migration

    One shard can accept the write while the other fails, producing divergence.

    11

    Ignoring Large-tenants and Hot Regions

    An apparently balanced assignment becomes uneven because traffic and data volume differ greatly among tenants or locations.

    12

    Skipping Resharding Tests

    Ownership cutover, rollback, cache invalidation, and in-flight write behaviour remain unverified.

    Recommended Test Cases

    Test Expected Evidence
    Point lookup A request with the shard key reaches one correct shard
    Missing shard key The application uses an approved lookup or rejects the request
    Range query The query reaches only the ranges required by its predicate
    Hash distribution Representative keys distribute within the approved balance objective
    Hot key One popular key does not destabilize unrelated shards
    Large tenant The tenant can be isolated or subdivided according to policy
    Directory failure Routing follows the documented cache and failure policy
    Stale directory cache The router detects the placement-version mismatch and refreshes safely
    Geo-routing attempt Changing an untrusted request parameter does not change authorized ownership
    Scatter-gather query Fan-out, merge, timeout, and partial failure remain bounded
    Cross-shard transaction Partial failure follows the documented coordination or compensation policy
    Shard split Source data and concurrent updates appear correctly on the new owners
    Shard move rollback Routing can safely return to the original owner
    Shard failure Only the intended ownership scope is affected where isolation permits
    Global uniqueness Two shards cannot independently violate the required invariant

    Sharding Best Practices

    Recommended Practices

    • Prove the storage or write bottleneck before introducing sharding.
    • Choose the shard key from real query and transaction patterns.
    • Measure both storage distribution and traffic distribution.
    • Keep related transactional records on the same shard where practical.
    • Prefer single-shard operations for routine request paths.
    • Use range sharding when ordered locality is a primary requirement.
    • Control sequential-key hotspots in range-sharded systems.
    • Use hash sharding when broad point-lookup distribution is more important than range locality.
    • Use consistent placement mechanisms when shard membership changes.
    • Use directory sharding when tenants require flexible individual placement.
    • Make the shard directory highly available, versioned, and recoverable.
    • Use geo sharding only with trusted regional ownership.
    • Document regional residency, replication, backup, and failover behaviour.
    • Use projections or analytics systems for frequent global queries.
    • Avoid routine scatter-gather transactions and reads.
    • Define global uniqueness separately from local shard constraints.
    • Replicate every shard according to its availability requirements.
    • Design online resharding and rollback before capacity is exhausted.
    • Monitor hot shards, hot tenants, hot keys, and mapping failures.
    • Test shard splits, moves, failures, and cross-shard workflows.

    Practice Exercise

    Design a sharding strategy for your online learning platform.

    Requirements

    1. Identify the database capacity limit that requires sharding.
    2. List the most frequent point, range, and global queries.
    3. Identify the most frequent transactions.
    4. Evaluate tenant ID as a shard key.
    5. Model small, medium, and very large tenants.
    6. Compare range, hash, directory, and geo placement.
    7. Keep enrollment and progress transactions local where practical.
    8. Design a global course-search projection.
    9. Define how email-to-account lookup finds the correct shard.
    10. Define global uniqueness requirements.
    11. Replicate each shard for availability.
    12. Design a tenant-move process with rollback.
    13. Define behaviour during stale directory routing.
    14. Test one shard failure.
    15. Monitor data, traffic, latency, and hot keys by shard.

    Sharding-design Template

    Data Shard Key Strategy Main Risk
    Tenant transactional data Tenant ID Directory or hash Large or busy tenant hotspot
    Activity events Tenant and time or distributed event key Hybrid hash and range Current-range write hotspot
    Global course search Search-document ID Separate search projection Projection freshness and authorization filtering
    Regional learner records Trusted tenant region and tenant ID Geo followed by directory or hash Regional imbalance and controlled tenant migration
    Global account lookup Normalized account lookup key Global directory or lookup index Index consistency and uniqueness

    Frequently Asked Questions

    1

    What is sharding?

    Sharding splits one logical dataset into subsets stored and processed by different physical database shards.

    2

    What is a shard key?

    A shard key is the value used to determine which shard owns a record or request.

    3

    What is range sharding?

    Range sharding assigns contiguous intervals of an ordered keyspace to different shards.

    4

    What is hash sharding?

    Hash sharding applies a deterministic hash to the shard key and maps the result to a shard or token range.

    5

    What is directory sharding?

    Directory sharding uses an explicit mapping that identifies the shard assigned to each tenant, key, or ownership group.

    6

    What is geo sharding?

    Geo sharding assigns data to regional shard groups according to trusted geographic, residency, or jurisdictional ownership.

    7

    Which strategy provides the best range queries?

    Range sharding normally provides the strongest locality when the query range uses the same key that defines shard boundaries.

    8

    Does hash sharding eliminate hot shards?

    It can distribute many keys evenly, but one high-traffic key or tenant can still overload its assigned shard.

    9

    What is scatter-gather?

    Scatter-gather sends a query to several shards and combines their partial results at a coordinating layer.

    10

    What is resharding?

    Resharding changes data ownership by splitting, moving, merging, or redistributing shard ranges or assignments.

    11

    Does sharding replace replication?

    No. Sharding distributes different records, while replication maintains additional copies of each shard for availability or reads.

    12

    Which sharding strategy should I choose?

    Choose according to query locality, traffic distribution, tenant size, transaction boundaries, rebalancing needs, geographic ownership, and the operational complexity the system can support.

    Key Takeaway

    Sharding divides one logical dataset across several physical data owners. Range sharding preserves ordered locality and supports efficient range scans, but sequential or uneven ranges can create hot shards. Hash sharding spreads many keys broadly and supports efficient point lookup, but loses natural range locality and requires careful resharding. Directory sharding provides flexible tenant-by-tenant placement, but the directory becomes critical routing infrastructure. Geo sharding places data according to trusted regional ownership, improving locality and supporting documented residency requirements while complicating global queries, migration, and failover. Choose the shard key from real access and transaction patterns, keep related data together, minimize scatter-gather and cross-shard transactions, replicate each shard, and design resharding before capacity is exhausted. Finally, monitor storage, traffic, latency, hot keys, hot tenants, routing failures, and migration progress independently for every shard.