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.

Advanced125–145 minutesShard-key locality/routing simulationSystem/composite/user/directory methodsDirect routing is connection-time shard selectionLast reviewed: August 2026

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.

01

Choose a sharding key from transaction locality, cardinality, growth and hotspot distribution rather than convenience.

02

Compare system-managed, composite, user-defined and 26ai directory-based distribution.

03

Explain sharding key versus super sharding key and the role of duplicated reference data.

04

Simulate direct key routing versus coordinator fan-out and quantify locality from ServiceHub operations.

05

Demonstrate a poor region-only key that creates imbalance/hotspots, then repair it with a locality-preserving higher-cardinality key.

Generation-time baseline, branding, licensing, topology, and tooling boundary

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.

sql · real topology schema shape
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.

sql · conceptual composite root table
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.

java · JDBC/UCP conceptual direct-routing shape
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

text · GDSCTL placement lookup on a deployed system
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

sql · setup
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

sql · requests by key and simulated shard
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

sql · two values cannot balance three/future shards well
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

sql · locality score for API inventory
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

sql · 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

  1. What business question should a sharding key answer?
  2. What additional key does composite distribution introduce?
  3. Why is a low-cardinality region-only key often poor?
  4. What happens when a sharding-aware client passes a shard key during connection checkout?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.