resharding
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 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
- Identify the migration ID and current phase.
- Confirm the source and target ownership metadata.
- Check historical-copy checkpoint and throughput.
- Check source change-generation and target apply rates.
- Check change-capture connection and retention.
- Check source and target CPU, memory, I/O, storage, and connections.
- Check failed and repeatedly retried batches.
- Check version-aware upsert behaviour.
- Check validation mismatches by key range.
- Check router placement versions and stale caches.
- Check queued work carrying old routing information.
- Check secondary indexes, caches, and projections.
- Pause or reduce migration load when production capacity is threatened.
- Resume, roll back, or complete through the approved migration procedure.
Common Resharding Mistakes
Waiting until the Shard Is Completely Full
Migration has no capacity headroom and competes with an already saturated production workload.
Copying Data without Capturing Concurrent Writes
Updates made after a record is copied can be missing from the destination.
Starting Change Capture after the Historical Copy
Writes occurring between the copy boundary and capture startup can be lost.
Using Unsafe Dual Writes
One shard can accept a write while the other shard fails, causing divergence.
Cutting Over before the Target Is Current
Requests reach a destination that is missing recent records or changes.
Validating Row Counts Only
The target can contain the correct number of rows but incorrect values, relationships, or business totals.
Ignoring Stale Routing Caches
Application instances continue sending requests to the old owner after cutover.
Leaving Two Write Authorities Active
Source and target can independently accept conflicting writes.
Forgetting Queued and Scheduled Work
Delayed jobs can execute using an ownership mapping captured before cutover.
Overloading Production with the Copy
Migration consumes source reads, target writes, network bandwidth, and connections required by user traffic.
Deleting the Source Too Early
Rollback and discrepancy investigation become impossible before the target has proved stable.
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
- Select one large tenant currently stored on a shared shard.
- Create a dedicated destination shard.
- Prepare the destination schema, indexes, replicas, and backups.
- Create a versioned migration record.
- Capture the migration-start position.
- Copy the tenant's historical records in bounded batches.
- Capture and replay concurrent changes.
- Make replay operations version-aware and idempotent.
- Validate counts, relationships, and business totals.
- Shadow-read the destination for selected queries.
- Reduce synchronization lag to the approved cutover boundary.
- Update the tenant-to-shard directory atomically.
- Refresh application routing caches.
- Monitor post-cutover errors and latency.
- Test rollback after target-only writes.
- 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
What is resharding?
Resharding changes how data is divided or assigned among database shards.
Why is resharding required?
It is required when existing shards become too large, uneven, overloaded, geographically unsuitable, or unable to support future growth.
What is a shard split?
A shard split divides one ownership range or group into two or more smaller shards.
What is a shard move?
A shard move transfers an existing ownership unit to another physical shard, server, or regional group.
What is online resharding?
Online resharding moves data while the application continues accepting reads and writes.
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.
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.
How is migrated data validated?
Compare record counts, checksums, key ranges, relationships, indexes, versions, and business-level totals at a known consistency boundary.
What is cutover?
Cutover transfers authoritative routing and write ownership from the source shard to the destination shard.
Why are placement versions needed?
Placement versions allow routers and shards to identify stale ownership metadata after a migration.
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.
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.