Table of Contents

    resharding

    DISTRIBUTED DATA & SHARDING

    Resharding

    Learn how resharding redistributes data when shards become too large, uneven, overloaded, or geographically unsuitable, and how shard splits, merges, tenant moves, consistent hashing, historical copying, change capture, dual reads, cutover, validation, rollback, replication, stale routing, and capacity controls support a safe migration.

    Introduction

    A sharded system divides one logical dataset across several physical database shards. The initial placement can work well when the system is small, but workload distribution changes over time.

    A shard can become problematic because:

    • Its storage approaches the available capacity
    • Its traffic grows faster than traffic on other shards
    • One tenant becomes substantially larger than expected
    • A key range receives most new writes
    • One geographic region receives growing demand
    • New database servers become available
    • Data must move to a different regional location
    • A previous shard-key or placement decision no longer fits the workload
    Before:
    
    Shard A = 20% of traffic
    Shard B = 65% of traffic
    Shard C = 15% of traffic
    
    
    Problem:
    
    Shard B is overloaded while
    other shards have spare capacity.

    Resharding changes data ownership so that selected records, tenants, key ranges, hash ranges, or geographic partitions move to different shards.

    Core idea: Resharding is not simply copying rows to another database. It is a controlled ownership transfer that must preserve correctness while data, concurrent writes, routing metadata, indexes, replicas, and application traffic move from an old shard layout to a new layout.

    Prerequisites

    # Prerequisite Why It Is Needed
    1 Range, hash, directory, and geo sharding The placement strategy determines which ownership units can move.
    2 Consistent hashing Hash-based placement can limit the key ranges affected by membership changes.
    3 Replication Every resulting shard normally requires replicas and validated failover.
    4 Change Data Capture Changes made during the historical copy must reach the destination.
    5 Transactions and idempotency Retries and partial writes must not cause duplicate or missing business operations.
    6 Load balancing and routing Applications need current ownership metadata during and after cutover.
    7 Observability Copy progress, change lag, validation, capacity, and routing failures require monitoring.

    What Is Resharding?

    Resharding is the process of changing how data is divided and assigned among database shards.

    Old ownership:
    
    Shard A owns Keys 1 to 1,000,000
    
    
    New ownership:
    
    Shard A owns Keys 1 to 500,000
    
    Shard B owns Keys 500,001 to 1,000,000

    Resharding can change:

    • The number of shards
    • The boundaries between key ranges
    • The owner of a hash or token range
    • The shard assigned to a tenant
    • The regional location of an ownership group
    • The physical servers hosting a shard
    • The shard key or partitioning strategy
    Resharding Flow
    plan new ownership → prepare destination → copy historical data → capture concurrent changes → validate → switch routing → retire old ownership

    Resharding vs Rebalancing

    Concept Meaning
    Resharding Changes logical data ownership, shard boundaries, shard count, or placement.
    Rebalancing Redistributes ownership or workload to improve capacity balance.
    Replica movement Moves an additional copy while logical primary ownership can remain unchanged.
    Node replacement Replaces infrastructure hosting an existing shard or replica.
    Migration A broader term covering data movement between systems, engines, regions, or schemas.

    Rebalancing can use resharding, but not every rebalancing action changes the logical shard structure.

    Why Reshard?

    Trigger Possible Resharding Action
    Shard storage is approaching capacity Split the shard or move selected ownership units
    One shard receives too much traffic Distribute hot ranges, keys, or tenants
    One tenant becomes very large Move the tenant to a dedicated shard or subdivide its data
    New database nodes are added Transfer ownership to use the additional capacity
    Several shards are underused Merge compatible shards or ownership ranges
    Regional ownership changes Move the tenant or partition to another regional shard group
    Shard key causes excessive fan-out Repartition data using a more suitable key

    Main Resharding Operations

    Shard Split

    A shard split divides one ownership unit into two or more smaller units.

    Before:
    
    Shard A owns Range 1 to 1,000
    
    
    After:
    
    Shard A owns Range 1 to 500
    
    Shard B owns Range 501 to 1,000

    A split is useful when a shard becomes too large or receives excessive traffic.

    Shard Merge

    A shard merge combines small compatible ownership units.

    Before:
    
    Shard A owns Range 1 to 100
    
    Shard B owns Range 101 to 200
    
    
    After:
    
    Shard A owns Range 1 to 200

    A merge can reduce operational overhead when several shards remain significantly underused.

    Shard Move

    A shard move transfers an existing ownership unit to another physical shard or server without necessarily changing its logical boundary.

    Before:
    
    Tenant 17 -> Shard A
    
    
    After:
    
    Tenant 17 -> Shard D

    Repartitioning

    Repartitioning changes the key or strategy used to distribute data.

    Old strategy:
    
    Range by sequential account ID
    
    
    New strategy:
    
    Hash by tenant ID

    Repartitioning is usually more complex than moving an existing ownership range because most or all records can receive a new destination.

    Strategy-specific Resharding

    Strategy Resharding Unit Main Concern
    Range sharding Contiguous key range Boundary selection and current-range hotspots
    Hash sharding Hash or token range Avoid remapping most keys when membership changes
    Directory sharding Tenant, account, or explicitly mapped group Atomically coordinating directory ownership with data movement
    Geo sharding Regional ownership group Residency, latency, backup, and failover requirements

    Range-shard Split

    Current boundary:
    
    Shard A owns 1 to 10,000,000
    
    
    New split point:
    
    5,000,000
    
    
    Target layout:
    
    Shard A owns 1 to 5,000,000
    
    Shard B owns 5,000,001 to 10,000,000

    The split point should consider both stored data and traffic. Dividing the numeric keyspace into equal intervals does not guarantee equal storage or workload.

    Consistent-hash Resharding

    Consistent hashing limits ownership movement when nodes join or leave.

    Before:
    
    Node A owns Token Range X
    
    
    New Node B joins inside Range X
          |
          v
    Node B receives part of Range X
    
    
    Only keys within the changed
    ownership interval need to move.

    Consistent hashing identifies affected ranges. The system still needs a safe data-copy, replication, validation, and routing-cutover process.

    Directory-based Tenant Move

    Old placement:
    
    tenant-17 -> shard-a
    
    
    Migration state:
    
    tenant-17 -> moving
    source = shard-a
    target = shard-d
    
    
    New placement:
    
    tenant-17 -> shard-d

    Directory sharding permits selective movement of one tenant without changing the placement formula for every other tenant.

    The directory must not indicate that the target is active before the target contains a validated and current copy of the tenant's data.

    Geo Resharding

    Geo resharding moves trusted ownership between regional shard groups.

    Old ownership:
    
    Tenant 17 -> Region A
    
    
    New requirement:
    
    Tenant 17 -> Region B
    
    
    Migration must address:
    
    - Transactional records
    - Object storage
    - Search projections
    - Backups
    - Replicas
    - Queue and event state
    - Routing metadata

    Regional movement should follow approved data-residency, backup, encryption, retention, disaster-recovery, and deletion requirements.

    Geo rule: Regional ownership must come from trusted policy and account data. A client must not trigger cross-region data movement by changing a request parameter.

    Offline Resharding

    Offline resharding stops or restricts writes while ownership changes.

    Enter maintenance mode
          |
          v
    Stop affected writes
          |
          v
    Copy data
          |
          v
    Validate destination
          |
          v
    Switch routing
          |
          v
    Resume writes

    Advantages

    • Simpler consistency model
    • No concurrent-write stream to capture
    • Clear cutover boundary
    • Easier comparison between source and destination

    Limitations

    • Requires maintenance or write interruption
    • Large datasets can create unacceptable downtime
    • Read behaviour during migration still requires definition
    • Queued or retrying clients can create a surge after reopening

    Online Resharding

    Online resharding keeps the application available while data moves.

    Capture migration start position
          |
          v
    Copy historical data
          |
          v
    Continue accepting writes
          |
          v
    Capture concurrent changes
          |
          v
    Replay changes on target
          |
          v
    Reduce migration lag
          |
          v
    Validate and cut over

    Online resharding reduces service interruption but introduces concurrent ownership and synchronization complexity.

    Snapshot plus Change Capture

    A common migration direction combines a consistent historical copy with an ordered stream of later changes.

    Migration begins at Position P
          |
          +-- Copy source state at P
          |
          +-- Capture writes after P
                  |
                  v
            Apply copied data
                  |
                  v
            Replay later changes
                  |
                  v
            Target reaches source

    The copy and change stream must share an unambiguous boundary. Otherwise, a change can be missed or applied twice.

    Migration Phases

    Phase Purpose
    Planning Define source, destination, ownership unit, capacity, and rollback
    Preparation Create destination schema, indexes, replicas, credentials, and monitoring
    Historical copy Transfer the existing records
    Change synchronization Capture and apply writes occurring after the copy boundary
    Validation Verify counts, checksums, constraints, relationships, and business totals
    Cutover Transfer authoritative routing to the destination
    Observation Monitor errors, latency, correctness, and migration completion
    Retirement Remove the old ownership copy after rollback is no longer required

    Historical Data Copy

    The historical copy transfers existing records to the destination.

    Read source in bounded batches
          |
          v
    Transform only when required
          |
          v
    Write destination idempotently
          |
          v
    Record checkpoint
          |
          v
    Continue next batch

    Batching should consider:

    • Source read capacity
    • Destination write capacity
    • Replication-log growth
    • Transaction size
    • Lock duration
    • Network throughput
    • Retry and checkpoint behaviour

    Copy-duration Estimate

    A simplified duration estimate is:

    \[ CopyDuration \approx \frac{ TotalDataToMove }{ EffectiveCopyThroughput } \]

    Effective throughput is constrained by the slowest relevant stage, including source reads, network transfer, destination writes, index maintenance, and validation overhead.

    Change-capture Lag

    Let \(G\) represent the source change-generation rate and \(A\) represent the destination change-application rate.

    The destination can catch up when:

    \[ A > G \]

    Migration lag grows when:

    \[ G > A \]

    If the destination cannot apply changes faster than they are generated, online migration cannot converge without reducing writes, increasing apply capacity, or changing the migration approach.

    Dual Write vs Change Capture

    Area Application Dual Write Change Capture
    Write path Application writes source and destination Application writes source, and changes are captured afterward
    Application changes Usually required Can be isolated in migration infrastructure
    Partial failure One write can succeed while the other fails Change stream can retry from an ordered position
    Ordering Must be handled across two writes Can follow a database change log where supported
    Recovery Requires reconciliation of divergent outcomes Requires retained change history and checkpoint recovery

    Dual-write rule: Two ordinary writes do not become atomic merely because the application sends them together. Design for one side succeeding while the other side times out or fails.

    Idempotent Migration Writes

    Copy and change events can be retried after uncertain failures.

    Migration writes record
          |
          v
    Destination commits record
          |
          v
    Acknowledgment is lost
          |
          v
    Migration retries same record

    The destination operation should produce the same correct result after a retry.

    Useful controls include:

    • Stable primary keys
    • Upserts with version checks
    • Operation identifiers
    • Source log positions
    • Conditional updates
    • Deduplication records

    Version-aware Upsert

    INSERT INTO course_progress
    (
        tenant_id,
        learner_id,
        course_id,
        progress_percentage,
        source_version
    )
    VALUES
    (
        :tenant_id,
        :learner_id,
        :course_id,
        :progress_percentage,
        :source_version
    )
    ON CONFLICT
    (
        tenant_id,
        learner_id,
        course_id
    )
    DO UPDATE SET
        progress_percentage =
            EXCLUDED.progress_percentage,
    
        source_version =
            EXCLUDED.source_version
    WHERE
        course_progress.source_version
            < EXCLUDED.source_version;

    Exact syntax and concurrency guarantees depend on the selected database. The key principle is that an older replay must not overwrite a newer state.

    Read-routing during Migration

    Reads require a defined source of truth while data moves.

    Strategy Direction Main Risk
    Read source only Keep all reads on the source until cutover Destination problems can remain hidden until late validation
    Shadow read Serve source result and compare silently with destination Additional load and privacy-sensitive comparison logs
    Destination-first with source fallback Read target and fall back when data is missing Can hide incomplete migration and create inconsistent latency
    Percentage cutover Route a controlled subset of eligible reads to the target Users can observe different versions if synchronization is incomplete
    Key-range cutover Move complete ownership units one at a time Requires accurate ownership metadata

    Shadow Reads

    Client request
          |
          v
    Read source
          |
          v
    Return source result to client
          |
          v
    Optionally read destination
    in controlled shadow path
          |
          v
    Compare results
          |
          v
    Record safe discrepancy metrics

    Shadow reads can validate real query behaviour without making the destination authoritative. Bound their traffic and avoid writing confidential field values to ordinary comparison logs.

    Routing Cutover

    Cutover changes which shard is authoritative for the moved ownership unit.

    Before cutover:
    
    Tenant 17 -> Source Shard
    
    
    Cutover transaction:
    
    Placement Version 41
          |
          v
    Placement Version 42
    
    
    After cutover:
    
    Tenant 17 -> Target Shard

    The placement change should be versioned so routers can identify outdated metadata.

    Conceptual Placement Record

    {
      "tenantId": 17,
      "sourceShard": "shard-a",
      "targetShard": "shard-d",
      "placementVersion": 42,
      "status": "active-on-target"
    }

    Stale Routing Caches

    Application instances can retain old range or directory mappings after cutover.

    Current ownership:
    
    Tenant 17 -> Shard D
    
    
    Router cache:
    
    Tenant 17 -> Shard A
    
    
    Result:
    
    Request reaches old owner.

    Possible controls include:

    • Versioned routing records
    • Bounded cache lifetime
    • Placement invalidation events
    • Wrong-owner responses
    • Safe forwarding during a limited transition
    • One bounded metadata-refresh retry

    Temporary Request Forwarding

    Request reaches old owner
          |
          v
    Old owner detects
    new placement version
          |
          v
    Forward or redirect through
    an approved internal mechanism
          |
          v
    Router refreshes mapping

    Forwarding should be temporary. Permanent forwarding chains increase latency, complicate troubleshooting, and can form loops.

    Data Validation

    A successful copy command does not prove that the complete destination is correct.

    Validation can include:

    • Record counts by ownership range
    • Checksums for bounded key ranges
    • Minimum and maximum keys
    • Null and constraint checks
    • Parent-child relationship checks
    • Business totals
    • Index availability
    • Query-result comparisons
    • Latest source and target positions

    Conceptual Batch Checksum

    SELECT
        COUNT(*) AS record_count,
        MIN(record_id) AS minimum_id,
        MAX(record_id) AS maximum_id,
        SUM(amount) AS total_amount
    FROM enrollment_records
    WHERE tenant_id = :tenant_id;

    Business totals and trusted hash comparisons can detect discrepancies that a row count alone would miss.

    Source and Target Validation

    For every migration batch:
    
    Source record count
        =
    Target record count
    
    
    Source business totals
        =
    Target business totals
    
    
    Source latest version
        =
    Target latest applied version
    
    
    All required indexes
        =
    Available and valid

    Validation rule: Validate data at a known consistency boundary. Comparing a changing source with a target at another point in time can report false differences.

    Rollback

    A rollback returns authoritative routing to the previous owner after a cutover problem.

    Cutover to target
          |
          v
    Critical issue detected
          |
          v
    Stop or restrict target writes
          |
          v
    Reconcile target changes
          |
          v
    Restore source authority
          |
          v
    Update routing version
          |
          v
    Validate service recovery

    Rollback becomes difficult after the target accepts writes that never reached the source.

    Define before cutover:

    • How target-only writes return to the source
    • How conflicts are detected
    • How long the old copy remains available
    • Which routing version represents rollback
    • What conditions make rollback unsafe
    • Who or what system can authorize rollback

    Point of No Return

    The point of no return is the stage after which simple routing rollback is no longer sufficient.

    Before point of no return:
    
    Source remains current
    or can be synchronized safely
    
    
    After point of no return:
    
    Target-only writes, schema changes,
    or source retirement require
    a forward recovery procedure

    Identify this boundary explicitly in the migration runbook.

    Rate Limiting the Migration

    Resharding consumes the same resources used by production traffic.

    Control:

    • Copy concurrency
    • Batch size
    • Source read rate
    • Destination write rate
    • Network bandwidth
    • Change-replay concurrency
    • Index-building activity
    • Validation-query concurrency
    Production pressure increases
          |
          v
    Reduce migration throughput
          |
          v
    Preserve user-facing capacity
    
    
    Production pressure decreases
          |
          v
    Increase migration throughput
    within approved boundaries

    Overload Protection

    Migration traffic should not cause the source, target, network, or shared dependencies to fail.

    Use:

    • Bounded worker pools
    • Rate limits
    • Backpressure
    • Pause and resume checkpoints
    • Database connection budgets
    • Replication-log retention monitoring
    • Production-priority resource controls

    Indexes and Constraints

    Destination schemas, indexes, and constraints should be prepared before authoritative cutover.

    Review:

    • Primary keys
    • Unique indexes
    • Foreign-key behaviour
    • Check constraints
    • Partition-local indexes
    • Global lookup indexes
    • Generated columns
    • Triggers and stored procedures

    Disabling constraints can increase copy speed, but requires complete post-copy validation before activation.

    Related Data Colocation

    Moving a parent record without its dependent records can create distributed joins or broken ownership.

    Tenant ownership unit:
    
    - Tenant profile
    - Users
    - Courses
    - Enrollments
    - Progress
    - Configuration
    - Idempotency records
    - Related workflow state

    Define the complete aggregate that must move together.

    Secondary Systems

    Transactional rows are only one part of the ownership footprint.

    Resharding can affect:

    • Search indexes
    • Distributed caches
    • Object-storage keys
    • Analytics projections
    • Queues and delayed jobs
    • Outbox records
    • Session or authorization data
    • Backup and retention configuration

    Each secondary system needs a decision about rebuilding, invalidating, moving, or retaining its data.

    Queued Work during Resharding

    Job created before cutover:
    
    tenantId = 17
    sourceShard = shard-a
    
    
    Job executes after cutover:
    
    tenantId = 17
    currentShard = shard-d

    Workers should resolve current ownership at execution time or carry a versioned routing context with an approved stale-routing policy.

    Cross-Shard Transactions during Migration

    Migration can temporarily create two physical copies of one logical record. The copies must not be treated as two independent authoritative records.

    Maintain one documented write authority at each migration stage, or use a datastore-supported atomic migration mechanism.

    Authority rule: At every moment, the system must know which shard is authoritative for new writes. Two copies do not imply two independent owners.

    Conceptual Migration State

    resharding:
      migrationId: tenant-17-move-001
    
      ownership:
        keyType: tenant
        keyValue: tenant-17
    
      source:
        shard: shard-a
        placementVersion: 41
    
      target:
        shard: shard-d
        placementVersion: 42
    
      phase: synchronizing
    
      copy:
        checkpoint: approved-checkpoint
        historicalComplete: true
    
      changeCapture:
        sourcePosition: approved-source-position
        targetAppliedPosition: approved-target-position
    
      validation:
        rowCounts: passed
        businessTotals: passed
        requiredIndexes: passed
    
      cutover:
        authorized: false
    
      rollback:
        supported: true
        sourceRetention: approved-policy

    This is a conceptual representation. Use the selected platform's supported online migration and resharding capabilities where available.

    PHP Version-aware Router

    <?php
    
    declare(strict_types=1);
    
    final class ShardRouter
    {
        public function __construct(
            private ShardDirectory $directory,
            private ShardConnectionFactory $connections
        ) {
        }
    
        public function connectionForTenant(
            int $tenantId
        ): ShardConnection {
            $placement =
                $this->directory->findActivePlacement(
                    $tenantId
                );
    
            if ($placement === null) {
                throw new RuntimeException(
                    'No active shard placement exists.'
                );
            }
    
            return $this->connections->create(
                shardId:
                    $placement->getShardId(),
    
                placementVersion:
                    $placement->getVersion()
            );
        }
    
        public function refreshAfterWrongOwner(
            int $tenantId,
            int $observedVersion
        ): ShardConnection {
            $this->directory->invalidateCachedPlacement(
                tenantId:
                    $tenantId,
    
                observedVersion:
                    $observedVersion
            );
    
            return $this->connectionForTenant(
                $tenantId
            );
        }
    }

    Production routing requires trusted tenant context, bounded retries, authenticated metadata, secure shard credentials, transaction boundaries, observability, and protection against cross-tenant routing.

    Security and Tenant Isolation

    Resharding must preserve security controls throughout the migration.

    • Use trusted tenant identity for data selection
    • Prevent clients from selecting source or destination shards
    • Encrypt migration traffic according to policy
    • Use least-privilege migration identities
    • Protect source and destination credentials
    • Preserve row-level and tenant-level access rules
    • Audit ownership and routing changes
    • Delete temporary copies according to approved retention
    • Avoid logging sensitive migrated values

    Backups and Recovery

    Resharding changes the data layout but does not remove backup requirements.

    Verify:

    • The source has a recoverable backup before migration
    • The target is included in backup policies
    • Point-in-time recovery boundaries are understood
    • A restore can reconstruct shard-directory ownership
    • Backups follow regional placement requirements
    • Source copies are not deleted before recovery validation

    Observability

    Useful resharding metrics include:

    • Total records and bytes to move
    • Copied records and bytes
    • Copy throughput
    • Estimated remaining ownership ranges
    • Change-generation rate
    • Change-application rate
    • Change-capture lag
    • Source and target CPU, memory, connections, and I/O
    • Validation mismatch count
    • Routing version by application instance
    • Wrong-owner response rate
    • Dual-read comparison mismatches
    • Migration retry count
    • Cutover and rollback events
    • Post-cutover latency and error rate

    Structured Migration Event

    {
      "migrationId": "tenant-17-move-001",
      "phase": "validation",
      "sourceShard": "shard-a",
      "targetShard": "shard-d",
      "placementVersion": 42,
      "historicalCopyComplete": true,
      "changeLag": 0,
      "validationResult": "passed"
    }

    Migration logs should not contain credentials, full connection strings, private payloads, or unprotected personal data.

    Alert Conditions

    Alert when:

    • Historical copy stops progressing
    • Change-capture lag continues increasing
    • The target apply rate remains below the source change rate
    • Source storage grows because change logs cannot be removed
    • Validation finds missing or conflicting records
    • Source or target production latency increases
    • Migration connections approach database limits
    • Routers continue using an old placement version
    • Wrong-owner responses increase after cutover
    • Target replicas are unhealthy
    • Cutover produces increased errors or latency
    • Rollback conditions are triggered

    Troubleshooting Workflow

    1. Identify the migration ID and current phase.
    2. Confirm the source and target ownership metadata.
    3. Check historical-copy checkpoint and throughput.
    4. Check source change-generation and target apply rates.
    5. Check change-capture connection and retention.
    6. Check source and target CPU, memory, I/O, storage, and connections.
    7. Check failed and repeatedly retried batches.
    8. Check version-aware upsert behaviour.
    9. Check validation mismatches by key range.
    10. Check router placement versions and stale caches.
    11. Check queued work carrying old routing information.
    12. Check secondary indexes, caches, and projections.
    13. Pause or reduce migration load when production capacity is threatened.
    14. Resume, roll back, or complete through the approved migration procedure.

    Common Resharding Mistakes

    1

    Waiting until the Shard Is Completely Full

    Migration has no capacity headroom and competes with an already saturated production workload.

    2

    Copying Data without Capturing Concurrent Writes

    Updates made after a record is copied can be missing from the destination.

    3

    Starting Change Capture after the Historical Copy

    Writes occurring between the copy boundary and capture startup can be lost.

    4

    Using Unsafe Dual Writes

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

    5

    Cutting Over before the Target Is Current

    Requests reach a destination that is missing recent records or changes.

    6

    Validating Row Counts Only

    The target can contain the correct number of rows but incorrect values, relationships, or business totals.

    7

    Ignoring Stale Routing Caches

    Application instances continue sending requests to the old owner after cutover.

    8

    Leaving Two Write Authorities Active

    Source and target can independently accept conflicting writes.

    9

    Forgetting Queued and Scheduled Work

    Delayed jobs can execute using an ownership mapping captured before cutover.

    10

    Overloading Production with the Copy

    Migration consumes source reads, target writes, network bandwidth, and connections required by user traffic.

    11

    Deleting the Source Too Early

    Rollback and discrepancy investigation become impossible before the target has proved stable.

    12

    Testing Cutover without Testing Rollback

    The procedure cannot safely recover when the target fails after receiving authoritative writes.

    Recommended Test Cases

    Test Expected Evidence
    Historical copy All records in the selected ownership unit reach the target
    Concurrent update An update made during the copy reaches the target in the correct order
    Repeated migration event Idempotency prevents duplicate or older state
    Deleted source record The deletion reaches the target according to the migration contract
    Copy-worker failure The worker resumes from a safe checkpoint
    Source write surge Change replay eventually catches up or migration pauses safely
    Target slowdown Backpressure protects the target and production workload
    Validation mismatch Cutover is blocked and the affected range is repaired
    Stale router The router refreshes placement and performs only a bounded retry
    Queued pre-cutover job The worker resolves current ownership before executing
    Cutover New traffic reaches the target without missing writes
    Post-cutover target failure The documented rollback or forward-recovery process succeeds
    Source retirement Removal occurs only after the observation and rollback period
    Tenant isolation The migration includes only the authorized ownership unit
    Backup restoration The new shard layout and directory can be reconstructed

    Resharding Best Practices

    Recommended Practices

    • Begin resharding before a shard reaches critical capacity.
    • Define the exact ownership unit being moved.
    • Measure storage and traffic distribution before selecting a split.
    • Prepare target schema, indexes, replicas, backups, and monitoring first.
    • Capture a clear historical-copy and change-stream boundary.
    • Make copy and replay operations idempotent.
    • Prevent older events from overwriting newer target data.
    • Keep one authoritative write owner at every migration stage.
    • Use bounded batches, connections, and migration concurrency.
    • Protect production traffic with backpressure and rate limits.
    • Validate row counts, checksums, relationships, and business totals.
    • Use shadow reads before authoritative cutover where appropriate.
    • Version shard-directory and ownership metadata.
    • Invalidate or refresh stale routing caches.
    • Resolve current ownership when delayed jobs execute.
    • Include caches, indexes, object storage, queues, and analytics projections.
    • Define rollback before cutover.
    • Retain the old copy until the target is proven stable.
    • Audit tenant and regional ownership changes.
    • Test copy, synchronization, cutover, failure, rollback, and recovery.

    Practice Exercise

    Design an online tenant resharding process for your learning platform.

    Requirements

    1. Select one large tenant currently stored on a shared shard.
    2. Create a dedicated destination shard.
    3. Prepare the destination schema, indexes, replicas, and backups.
    4. Create a versioned migration record.
    5. Capture the migration-start position.
    6. Copy the tenant's historical records in bounded batches.
    7. Capture and replay concurrent changes.
    8. Make replay operations version-aware and idempotent.
    9. Validate counts, relationships, and business totals.
    10. Shadow-read the destination for selected queries.
    11. Reduce synchronization lag to the approved cutover boundary.
    12. Update the tenant-to-shard directory atomically.
    13. Refresh application routing caches.
    14. Monitor post-cutover errors and latency.
    15. Test rollback after target-only writes.
    16. Retire the old copy only after the approved observation period.

    Resharding-design Template

    Decision Selected Direction Primary Risk Controlled
    Ownership unit Tenant, key range, or hash range Prevents incomplete or unrelated data movement
    Historical state Checkpointed bounded copy Supports restart after migration-worker failure
    Concurrent changes Ordered Change Data Capture stream Prevents updates from being lost during copying
    Write authority One active owner per migration phase Prevents source-target divergence
    Validation Counts, checksums, relationships, and business totals Detects incomplete or incorrect target data
    Routing Versioned atomic placement update Prevents applications from using incompatible ownership maps
    Rollback Retained source with target-write reconciliation Supports recovery from post-cutover failure
    Capacity protection Bounded copy and replay throughput Prevents migration from overloading production

    Frequently Asked Questions

    1

    What is resharding?

    Resharding changes how data is divided or assigned among database shards.

    2

    Why is resharding required?

    It is required when existing shards become too large, uneven, overloaded, geographically unsuitable, or unable to support future growth.

    3

    What is a shard split?

    A shard split divides one ownership range or group into two or more smaller shards.

    4

    What is a shard move?

    A shard move transfers an existing ownership unit to another physical shard, server, or regional group.

    5

    What is online resharding?

    Online resharding moves data while the application continues accepting reads and writes.

    6

    How are writes handled during an online move?

    A common direction is to copy historical data from a known position and capture later changes until the destination catches up.

    7

    Why are ordinary dual writes risky?

    The source write can succeed while the destination write fails, or the writes can be applied in different orders.

    8

    How is migrated data validated?

    Compare record counts, checksums, key ranges, relationships, indexes, versions, and business-level totals at a known consistency boundary.

    9

    What is cutover?

    Cutover transfers authoritative routing and write ownership from the source shard to the destination shard.

    10

    Why are placement versions needed?

    Placement versions allow routers and shards to identify stale ownership metadata after a migration.

    11

    When can the old shard copy be deleted?

    Delete it only after target validation, post-cutover observation, backup verification, and the approved rollback period have completed.

    12

    Does consistent hashing perform resharding automatically?

    Consistent hashing identifies changed ownership ranges and limits remapping. A separate migration process must still copy, synchronize, validate, and activate authoritative data safely.

    Key Takeaway

    Resharding changes ownership when the current shard layout no longer provides balanced capacity, storage, locality, or operational safety. Common operations include shard splits, merges, tenant moves, hash-range transfers, and regional ownership changes. Safe online resharding begins with a known copy boundary, transfers historical records, captures concurrent changes, replays those changes idempotently, validates the destination, and switches routing through versioned ownership metadata. Keep one authoritative write owner at every phase, because two copies must not become two independent authorities. Protect production capacity with bounded migration throughput, detect stale routing caches, include delayed jobs and secondary systems, and retain the source until rollback is no longer required. Finally, treat cutover and rollback as equally important procedures and test failures during copying, synchronization, validation, routing, and post-cutover operation.