Chapter 26 · Sharding, Distributed Database Concepts, and Global Data Placement
Sharding Keys, Data Distribution, Duplicated Tables, Locality, and Query Routing
Choose shard and super-shard keys from transaction locality, growth and hotspot evidence; compare system-managed, composite, user-defined and directory-based placement and show how key-based connection routing differs from coordinator fan-out.
Learning outcomes
ServiceHub can distribute customers evenly yet still perform poorly if every work-order request needs data from two regions. The sharding key is therefore a transaction-locality design choice, not merely a hash column. It controls row placement and can also be passed during connection checkout so the driver/shard director sends the session directly to the shard that owns the tenant's data.
Choose a sharding key from transaction locality, cardinality, growth and hotspot distribution rather than convenience.
Compare system-managed, composite, user-defined and 26ai directory-based distribution.
Explain sharding key versus super sharding key and the role of duplicated reference data.
Simulate direct key routing versus coordinator fan-out and quantify locality from ServiceHub operations.
Demonstrate a poor region-only key that creates imbalance/hotspots, then repair it with a locality-preserving higher-cardinality key.
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 good key answers “what should be together?”
If almost every ServiceHub transaction reads/updates one
tenant's technicians, work orders and notes,
TENANT_ID is a natural root sharding-key candidate.
Child tables should carry/equi-partition on the same key so one
tenant transaction stays local. A date or status column may
spread rows but scatter the business aggregate that applications
update together.
2. System-managed distribution
System-managed distribution uses consistent hashing. Applications do not map key ranges to specific shards; Oracle distributes chunks evenly and automatically rebalances them when shards are added or removed. It is a good fit when you want uniform distribution and do not need explicit geographic ownership of individual key ranges.
CREATE SHARDED TABLE servicehub_tenant ( tenant_id NUMBER NOT NULL, tenant_name VARCHAR2(100) NOT NULL, CONSTRAINT servicehub_tenant_pk PRIMARY KEY(tenant_id))PARTITION BY CONSISTENT HASH (tenant_id)PARTITIONS AUTOTABLESPACE SET servicehub_ts;
3. Composite distribution: geography first, balance second
Composite sharding first places data by a super sharding key into a shardspace (for example legal region), then consistent-hashes the sharding key inside that shardspace. This is useful when geography/data sovereignty must be explicit while tenants within the region still need balanced distribution.
CREATE SHARDED TABLE servicehub_tenant ( legal_region VARCHAR2(20) NOT NULL, tenant_id NUMBER NOT NULL, tenant_name VARCHAR2(100) NOT NULL, CONSTRAINT servicehub_tenant_pk PRIMARY KEY(legal_region,tenant_id))PARTITIONSET BY LIST (legal_region)PARTITION BY CONSISTENT HASH (tenant_id)PARTITIONS AUTO( PARTITIONSET europe VALUES ('EU') TABLESPACE SET servicehub_eu_ts, PARTITIONSET middle_east VALUES ('ME') TABLESPACE SET servicehub_me_ts);
The exact partition-set values/tablespace sets are catalog design and are omitted here because the mandatory lab is not a deployed sharded database.
4. User-defined and directory-based placement
User-defined sharding explicitly maps user-defined ranges/lists/partitions to shards/shardspaces and is appropriate when existing placement has to be preserved. 26ai adds directory-based distribution, an enhanced explicit-placement model for workloads with too few distinct key values for even consistent-hash distribution or where administrators require direct key-to-shard control.
5. Duplicated tables preserve local joins
Small common reference data—status codes, equipment catalog, region definitions—can be duplicated on every shard. A tenant-local work-order transaction then joins local sharded rows to a local duplicated copy instead of reaching another database. The tradeoff is duplicated-table refresh/synchronization cost and catalog dependence for updates.
6. Direct routing selects the shard at connection checkout
Oracle sharding-aware clients can pass a sharding key during connection acquisition. The connection layer caches routing metadata and establishes the session directly to the shard owning that key. With composite sharding, both sharding key and super sharding key may be required. Once connected to that shard, ordinary SQL/DML runs within that shard unless the application instead uses coordinator/proxy-routing paths.
OracleShardingKey shardKey = pool.createShardingKeyBuilder() .subkey(tenantId, OracleType.NUMBER) .build();Connection c = pool.createConnectionBuilder() .shardingKey(shardKey) .build();
Use the current driver APIs/version supported by your application; this is the routing concept, not a reason to hard-code shard hostnames.
7. Real routing evidence
GDSCTL> config chunks -key 42017 -table_family SERVICEHUB_OWNER.SERVICEHUB_TENANT
config chunks -key tells operators which
chunk/shard owns a key. Application telemetry should also record
service/shard identity so a supposed single-shard request can be
verified.
8. Free simulation: map tenants to three logical shards
BEGIN EXECUTE IMMEDIATE 'DROP TABLE sh26_route_workload PURGE';EXCEPTION WHEN OTHERS THEN IF SQLCODE != -942 THEN RAISE; END IF; END;/CREATE TABLE sh26_route_workload ( request_id NUMBER PRIMARY KEY, tenant_id NUMBER NOT NULL, legal_region VARCHAR2(10) NOT NULL, operation_name VARCHAR2(30) NOT NULL, simulated_shard VARCHAR2(20) NOT NULL);INSERT INTO sh26_route_workloadSELECT LEVEL, CASE WHEN LEVEL <= 70 THEN 1 ELSE LEVEL END AS tenant_id, CASE WHEN MOD(LEVEL,2)=0 THEN 'EU' ELSE 'ME' END, 'WORK_ORDER_LOOKUP', CASE MOD( CASE WHEN LEVEL <= 70 THEN 1 ELSE LEVEL END ,3) WHEN 0 THEN 'SHARD_1' WHEN 1 THEN 'SHARD_2' ELSE 'SHARD_3' ENDFROM dualCONNECT BY LEVEL <= 100;COMMIT;
This intentionally creates a hot Tenant 1 generating 70% of the requests. Even perfect row distribution cannot remove a business hot key whose workload itself is concentrated.
9. Measure locality and hotspot pressure
SELECT tenant_id, simulated_shard, COUNT(*) AS request_countFROM sh26_route_workloadGROUP BY tenant_id,simulated_shardORDER BY request_count DESCFETCH FIRST 10 ROWS ONLY;SELECT simulated_shard,COUNT(*) AS request_countFROM sh26_route_workloadGROUP BY simulated_shardORDER BY request_count DESC;
Expected: whichever simulated shard owns Tenant 1 becomes hot despite a nominal three-shard design. Sharding cannot split one indivisible hot key unless the application/schema chooses a finer key or an additional hierarchy.
10. Deliberately wrong: shard only by low-cardinality region
SELECT legal_region, COUNT(*) AS requestsFROM sh26_route_workloadGROUP BY legal_regionORDER BY legal_region;
A region-only key gives two routing buckets, couples all tenants in one region to the same failure/capacity boundary and makes regional hotspots hard to rebalance. The repair is typically composite design: region as super key for sovereignty, tenant/customer as high-cardinality sharding key for distribution inside the region.
11. Cross-shard work begins when APIs lose the key
CREATE TABLE sh26_api_inventory ( api_name VARCHAR2(60) PRIMARY KEY, carries_tenant_key CHAR(1) CHECK (carries_tenant_key IN ('Y','N')), expected_scope VARCHAR2(20));INSERT INTO sh26_api_inventory VALUES('GET_WORK_ORDER','Y','SINGLE_SHARD');INSERT INTO sh26_api_inventory VALUES('UPDATE_WORK_ORDER','Y','SINGLE_SHARD');INSERT INTO sh26_api_inventory VALUES('GLOBAL_REVENUE_REPORT','N','MULTI_SHARD');INSERT INTO sh26_api_inventory VALUES('SEARCH_ALL_TENANTS','N','MULTI_SHARD');COMMIT;SELECT expected_scope,COUNT(*) AS apisFROM sh26_api_inventoryGROUP BY expected_scope;
The key should appear naturally in latency-sensitive transactional APIs. Global reporting remains legitimate, but it should be recognized as fan-out/coordinator work with different SLOs.
12. Cleanup
DROP TABLE sh26_api_inventory PURGE;DROP TABLE sh26_route_workload PURGE;
13. Production judgment
Choose the shard key from the dominant transaction boundary and growth distribution. Use system-managed for automatic balanced hashing, composite for explicit shardspace/sovereignty plus hashing, user-defined/directory-based when placement itself is business policy. Duplicate small reference tables to preserve local joins. Keep global/fan-out APIs explicit instead of hiding them behind the same latency contract as key-routed OLTP.
No mandatory lab deploys actual shards. On a real deployment, direct routing requires a sharding-aware client/connection pool and valid global service; never route by hard-coded host. Lesson 3 now investigates the requests that cannot stay local and why cross-shard queries/transactions create coordinator, network and failure-state costs.
Check your understanding
- What business question should a sharding key answer?
- What additional key does composite distribution introduce?
- Why is a low-cardinality region-only key often poor?
- What happens when a sharding-aware client passes a shard key during connection checkout?
- Why are duplicated tables useful to a tenant-local transaction?
Review the answers
Which rows/operations should be colocated so the common transaction stays on one shard.
A super sharding key that selects the shardspace/placement domain before hashing the shard key inside it.
It creates too few placement buckets and couples many tenants to one hotspot/failure domain.
The routing layer connects directly to the shard owning that key, avoiding catalog-coordinated fan-out for that session.
They keep small shared reference data local so joins do not require cross-shard access.
Authoritative references
- Data Distribution Methods — system/composite/user/directory methods
- Direct Routing to a Shard — sharding-key connection routing
- Config Chunks — key/chunk placement evidence
- Creating Duplicated Tables — duplicated reference data
- Changes in 26ai — directory distribution/synchronous duplicated tables