Chapter 26 · Sharding, Distributed Database Concepts, and Global Data Placement
Resharding, Scaling, Deployment Topology, Data Sovereignty, and Operational Complexity
Treat add-shard, chunk movement, split/move partition-set operations, regional placement, replication and data sovereignty as planned operational migrations with monitoring, failure handling and rollback—not frictionless horizontal scaling.
Learning outcomes
ServiceHub's EU shardgroup approaches capacity while another region is lightly loaded. “Add a shard” is not the end of the operation: data must move, routing caches/topology must converge, backup/recovery must include the new failure domain, and residency rules must still be satisfied while chunks are in transit. Resharding is a controlled data-placement migration.
Explain chunks, add-shard/rebalance and move-chunk operations for system-managed placement.
Explain shardspace/shardgroup placement, replication roles and regional sovereignty constraints.
Understand 26ai composite split-partitionset/bulk-move and automatic sharding-key row movement enhancements.
Model a shard expansion with prechecks, move states, failure/rollback and post-move verification.
Connect resharding to OMF/TDE/backup/upgrade/network prerequisites instead of treating it as a single GDSCTL command.
Mandatory examples target Oracle AI Database Free 26ai and were reviewed against the current August 2026 licensing documentation, RU 23.26.3, SQL Developer 26.2, and SQLcl 26.2.1.222.1617. Oracle documentation now uses Oracle Globally Distributed AI Database / Oracle Globally Distributed Database for the technology formerly called Oracle Sharding. The 26ai licensing matrix marks Oracle Globally Distributed Database available in Free with use limited to three shards. Raft replication is also marked available in Free with use limited to three nodes, and each Free node remains limited to 2 CPU cores, 2 GB RAM, and 12 GB user data. However, the traditional shard-catalog creation guide still documents the shard catalog as an Enterprise Edition PDB and explicitly disallows CDB$ROOT as the catalog. Therefore this chapter does not claim that a complete catalog + shard directors + multi-host shard topology is a zero-cost single-machine lab. Mandatory work uses one Free FREEPDB1 database to simulate placement, routing, fan-out and failure semantics while all GDSCTL/SHARD DDL examples are clearly labeled topology-dependent. Shard directors may run on separate servers without a separate shard-director server license. Full Data Guard, Active Data Guard, RAC, GoldenGate, Exadata and multi-host production topologies retain their own entitlement/operational requirements. No lab raises COMPATIBLE, changes hidden parameters, or modifies the GitHub repository.
1. A chunk is the practical movement unit
System-managed distribution hashes key ranges into many logical chunks. A shard owns many chunks. When capacity changes, Oracle moves chunks rather than rehashing every row to every shard. Consistent hashing plus many chunks lets a subset of data move when a shard is added/removed.
GDSCTL> config chunksGDSCTL> config chunks -show_reshardGDSCTL> config shardGDSCTL> config task
2. Add shard is a deployment plus data-movement workflow
GDSCTL> add shard -connect newshard-host:1521/newshardpdb -pwd <gsmuser-password> -shardgroup primary_baku -cdb NEWCDBGDSCTL> deployGDSCTL> config shard -shard NEWCDB_NEWSHARDPDBGDSCTL> config chunks -show_reshard
A shard can reach Deployed before background
rebalancing has fully moved its intended chunks. Production
readiness therefore requires both deployment state and
data-placement/task completion evidence.
3. OMF is part of chunk-movement mechanics
Current shard-database preparation documentation requires a
valid DB_CREATE_FILE_DEST so Oracle Managed Files
(OMF) can hold the transportable tablespaces used by chunk
movement/rebalancing. Character sets also need to match where
movement mechanisms require it.
SHOW PARAMETER compatibleSHOW PARAMETER db_create_file_destSELECT parameter,valueFROM nls_database_parametersWHERE parameter IN ( 'NLS_CHARACTERSET', 'NLS_NCHAR_CHARACTERSET')ORDER BY parameter;
4. TDE key consistency affects movement
When Transparent Data Encryption (TDE) protects sharded tablespaces, chunk movement needs compatible encryption-key access across shards. Current sharding guidance requires the participating shards to share/use the appropriate encryption key for encrypted tablespace movement. Key-management design therefore becomes part of resharding runbooks.
5. Composite layouts encode data sovereignty
Composite sharding can map a super key such as legal region to a shardspace, then hash tenant IDs within that shardspace. That allows “EU data stays in EU shardspace” while scaling across multiple EU shards. Adding capacity must preserve the shardspace placement rule; moving a chunk to the wrong shardspace would be a residency failure even if the database stayed available.
6. 26ai: split/move partition-set enhancements
26ai improves composite redistribution.
SPLIT PARTITIONSET can split existing chunks by
super-shard-key values into new shardspaces and bulk-move data,
keeping data online for as much of the operation as possible.
MOVE PARTITIONSET and
MODIFY PARTITIONSET add further placement
maintenance. These are 26ai features; do not backport them
conceptually to 21c operational runbooks.
7. 26ai: automatic data movement after sharding-key update
Updating a sharding key can now trigger Oracle-managed row movement within a shard or between shards rather than forcing the application to delete/reinsert manually. That convenience still changes the row's failure/latency domain and may move related table-family data; sharding-key updates should remain rare, intentional business events.
8. Free simulation: plan a three-to-four logical-shard expansion
BEGIN EXECUTE IMMEDIATE 'DROP TABLE sh26_chunk_move PURGE';EXCEPTION WHEN OTHERS THEN IF SQLCODE != -942 THEN RAISE; END IF; END;/BEGIN EXECUTE IMMEDIATE 'DROP TABLE sh26_chunk_map PURGE';EXCEPTION WHEN OTHERS THEN IF SQLCODE != -942 THEN RAISE; END IF; END;/CREATE TABLE sh26_chunk_map ( chunk_id NUMBER PRIMARY KEY, region_name VARCHAR2(10) NOT NULL, shard_name VARCHAR2(20) NOT NULL, row_estimate NUMBER NOT NULL);INSERT INTO sh26_chunk_mapSELECT LEVEL, CASE WHEN LEVEL <= 12 THEN 'EU' ELSE 'ME' END, CASE WHEN LEVEL BETWEEN 1 AND 6 THEN 'SHARD_EU_1' WHEN LEVEL BETWEEN 7 AND 12 THEN 'SHARD_EU_2' ELSE 'SHARD_ME_1' END, 1000 + MOD(LEVEL*137,500)FROM dualCONNECT BY LEVEL <= 18;CREATE TABLE sh26_chunk_move ( chunk_id NUMBER PRIMARY KEY, source_shard VARCHAR2(20) NOT NULL, target_shard VARCHAR2(20) NOT NULL, move_state VARCHAR2(20) NOT NULL, started_at TIMESTAMP, finished_at TIMESTAMP, CONSTRAINT sh26_move_state_ck CHECK (move_state IN ( 'PLANNED','COPYING','CUTOVER','DONE','FAILED' )));COMMIT;
9. Select EU chunks to rebalance—without violating residency
INSERT INTO sh26_chunk_move( chunk_id,source_shard,target_shard,move_state)SELECT chunk_id, shard_name, 'SHARD_EU_3', 'PLANNED'FROM sh26_chunk_mapWHERE region_name='EU' AND MOD(chunk_id,3)=0;COMMIT;SELECT m.chunk_id, c.region_name, m.source_shard, m.target_shard, m.move_stateFROM sh26_chunk_move mJOIN sh26_chunk_map c ON c.chunk_id=m.chunk_idORDER BY m.chunk_id;
Every planned target is an EU shard. The simulation makes residency an explicit validation rule rather than an operational assumption.
10. Deliberately wrong: move an EU chunk to the ME shard
SELECT c.chunk_id,c.region_name,m.target_shardFROM sh26_chunk_map cJOIN ( SELECT 3 AS chunk_id,'SHARD_ME_1' AS target_shard FROM dual) m ON m.chunk_id=c.chunk_idWHERE c.region_name='EU' AND m.target_shard LIKE 'SHARD_ME_%';-- Expected: chunk 3 appears as a residency-policy violation.
A distributed system can be technically healthy while violating residency policy. Placement governance must be a precondition/check, not a post-incident audit.
11. Simulate move and cutover states
UPDATE sh26_chunk_moveSET move_state='COPYING', started_at=SYSTIMESTAMPWHERE move_state='PLANNED';UPDATE sh26_chunk_moveSET move_state='CUTOVER'WHERE chunk_id=3;UPDATE sh26_chunk_moveSET move_state='DONE', finished_at=SYSTIMESTAMPWHERE chunk_id=3;UPDATE sh26_chunk_mapSET shard_name='SHARD_EU_3'WHERE chunk_id=3;COMMIT;SELECT * FROM sh26_chunk_move ORDER BY chunk_id;SELECT shard_name,SUM(row_estimate) estimated_rowsFROM sh26_chunk_mapGROUP BY shard_nameORDER BY shard_name;
Real Oracle resharding has its own transport/catch-up/routing state; these labels only teach operational checkpoints.
12. Failure/rollback runbook
- Preflight catalog/director/shard health and network reachability.
- Verify backup/recovery and enough target/source storage.
- Record chunk ownership and residency policy before move.
- Pause conflicting upgrade/patch/add-shard operations as documented.
-
Monitor
config chunks -show_reshard,config task, alert/ADR and catalog DDL/task failures. - On failure, do not manually delete copied tablespaces/rows; use supported suspend/cancel/recover commands and Oracle recovery guidance.
- After completion, verify key routing, row counts/checksums, duplicated-table state, backup coverage and application latency.
13. Upgrade sequencing interacts with movement
Current guidance says complete pending
MOVE CHUNK operations and do not start new chunk
moves/add shards during distributed database upgrade. Upgrade
order is catalog first, then shard directors, then shards.
Horizontal scaling and patching are therefore coordinated fleet
operations, not independent DBA tasks.
14. Cleanup
DROP TABLE sh26_chunk_move PURGE;DROP TABLE sh26_chunk_map PURGE;
15. Production judgment
Scale out because measured shard capacity/SLO requires it, not because “distributed systems scale horizontally.” Resharding consumes network/storage/redo/CPU, changes routing and failure placement, and must preserve sovereignty and encryption/backup guarantees. Keep capacity headroom so you can move data before a shard is already saturated.
The local simulation performs no real chunk movement. Real system-managed movement depends on deployed GDS topology/OMF and current operational restrictions. 26ai adds composite bulk split/move and automatic sharding-key row movement; check the exact RU/known issues before using them. Lesson 5 now compares whether sharding is the correct architecture at all versus RAC, Data Guard, partitioning or application-owned distribution.
Check your understanding
- What is the purpose of a sharding chunk?
- Why is DB_CREATE_FILE_DEST relevant to chunk movement?
- What does a composite shardspace let you encode?
- What operational conflict should be avoided during upgrades?
- Why is adding a shard not complete when it first shows Deployed?
Review the answers
It is a logical movement/placement unit so rebalancing can move subsets of data instead of rehashing everything.
OMF storage is used by the transportable-tablespace/chunk movement infrastructure.
Explicit placement such as legal region/data sovereignty, with hashing inside that domain.
Finish/pause chunk moves and do not add shards during the upgrade sequence.
Deployment/schema catch-up can finish before background chunk rebalancing/data movement completes.
Authoritative references
- Changes in 26ai — split/move/key-update enhancements
- Create the Shard Databases — COMPATIBLE/OMF/character-set prerequisites
- Config Chunks — reshard monitoring
- Using Transparent Data Encryption — TDE movement key requirements
- Patching and Upgrading — move/add-shard/upgrade sequencing