Chapter 26 · Sharding, Distributed Database Concepts, and Global Data Placement
Cross-Shard Queries/Transactions, Consistency, Failure Domains, and Application Design
Explain coordinator fan-out, global read consistency and two-phase commit for cross-shard DML; design idempotent locality-preserving APIs and rehearse partial-failure/in-doubt behavior without pretending distributed transactions behave like local ones.
Learning outcomes
Most ServiceHub OLTP is tenant-local, but finance requests a global revenue aggregation and an operations workflow updates records for two tenants in one business transaction. Those requests cannot be treated like a single local transaction. A multi-shard query is decomposed/fanned out to multiple shards and merged by a coordinator. A cross-shard DML transaction becomes a distributed transaction whose commit can require two-phase commit (2PC) and can enter an in-doubt state if communication fails between prepare and commit.
Explain query coordinator fan-out/partial aggregation and why single-shard SQL is a different latency class.
Explain global read consistency/SCN coordination for multi-shard reads and current consistency modes at architecture level.
Explain 26ai cross-shard DML and two-phase commit, including in-doubt transaction recovery.
Design locality-preserving, idempotent application APIs with explicit retry/request IDs.
Simulate partial shard failure and distributed transaction state on one Free database without pretending it is actual 2PC.
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. Multi-shard SELECT is distributed execution
When a query cannot be mapped to one shard, the shard
catalog/query coordinator transforms the statement into
shard-local work, sends it to relevant shards and may perform a
final merge/aggregation. For example, each shard can compute a
partial COUNT/SUM, while the coordinator combines
those partial results. Network transfer and slowest-shard
latency now matter.
SELECT tenant_id,SUM(amount)FROM servicehub_ordersGROUP BY tenant_id;
With no single shard key, this is fan-out work. Supported query
constructs have distributed restrictions; for example, current
sharding documentation lists unsupported constructs such as
CONNECT BY and MODEL in multi-shard
query processing.
2. Global consistency costs coordination
A distributed read needs a consistent logical point across participating shards when strict read consistency is requested. Oracle coordinates System Change Numbers (SCNs) for the query. Alternative locality/standby modes can trade freshness/availability/region affinity against strict synchronization. Do not promise “local read latency” for a query that requires global synchronized state.
3. 26ai parallel cross-shard DML exists—but physics remains
26ai adds parallel cross-shard DML support so the query coordinator can run updates/inserts/deletes against multiple shards concurrently rather than serially. That improves execution parallelism, but the operation is still distributed: multiple databases participate, network failures exist, and transaction commit semantics still require coordination.
UPDATE servicehub_ordersSET status_code='REVIEW'WHERE status_code='OPEN';-- Without a shard-key restriction this can affect multiple shards.-- The coordinator performs distributed DML and commit coordination.
4. Two-phase commit preserves atomicity across databases
In the prepare phase, the global coordinator asks participants whether they can commit. If all are prepared, the commit phase makes the outcome durable across participants. A crash/network break after prepare can leave a transaction in doubt: participating databases retain transaction state/locks until recovery resolves the global outcome.
SELECT local_tran_id, global_tran_id, state, mixed, advice, tran_comment, fail_time, retry_timeFROM dba_2pc_pendingORDER BY fail_time;SELECT *FROM dba_2pc_neighborsORDER BY local_tran_id;
These views are for genuine distributed transaction recovery evidence. The mandatory local simulation below does not populate them.
5. 26ai recovery protocol awareness
Oracle 26ai adds the
ALLOW_LEGACY_RECO_PROTOCOL control for distributed
transaction recovery. The default remains compatible with older
peers. If set to the upgraded protocol mode, every database
participating in distributed transactions must be 26ai or later;
otherwise recovery can fail. This is a migration/security
control, not a tuning knob.
SELECT name,value,issys_modifiable,ispdb_modifiableFROM v$parameterWHERE name='allow_legacy_reco_protocol';
6. Free simulation: partial aggregation across logical shards
BEGIN EXECUTE IMMEDIATE 'DROP TABLE sh26_xshard_orders PURGE';EXCEPTION WHEN OTHERS THEN IF SQLCODE != -942 THEN RAISE; END IF; END;/CREATE TABLE sh26_xshard_orders ( simulated_shard VARCHAR2(20) NOT NULL, tenant_id NUMBER NOT NULL, work_order_id NUMBER NOT NULL, amount NUMBER(10,2) NOT NULL, status_code VARCHAR2(12) NOT NULL, PRIMARY KEY(simulated_shard,tenant_id,work_order_id));INSERT INTO sh26_xshard_orders VALUES('SHARD_1',1,101,120,'OPEN');INSERT INTO sh26_xshard_orders VALUES('SHARD_1',4,104,150,'OPEN');INSERT INTO sh26_xshard_orders VALUES('SHARD_2',2,102,200,'OPEN');INSERT INTO sh26_xshard_orders VALUES('SHARD_3',3,103,300,'CLOSED');INSERT INTO sh26_xshard_orders VALUES('SHARD_3',6,106,175,'OPEN');COMMIT;
SELECT simulated_shard,SUM(amount) AS shard_sumFROM sh26_xshard_ordersGROUP BY simulated_shardORDER BY simulated_shard;SELECT SUM(shard_sum) AS global_sumFROM ( SELECT simulated_shard,SUM(amount) AS shard_sum FROM sh26_xshard_orders GROUP BY simulated_shard);
The second statement mimics coordinator final aggregation conceptually. It runs in one database and therefore has none of the real network/slow-shard/SCN coordination cost.
7. Simulate distributed work state and failure
CREATE TABLE sh26_distributed_ops ( request_id VARCHAR2(40) PRIMARY KEY, shard_name VARCHAR2(20) NOT NULL, prepare_state VARCHAR2(12) NOT NULL, commit_state VARCHAR2(12) NOT NULL, CONSTRAINT sh26_prepare_ck CHECK (prepare_state IN ('NEW','PREPARED','FAILED')), CONSTRAINT sh26_commit_ck CHECK (commit_state IN ('PENDING','COMMITTED','ROLLED_BACK')));INSERT INTO sh26_distributed_opsVALUES('REQ-9001','SHARD_1','PREPARED','PENDING');INSERT INTO sh26_distributed_opsVALUES('REQ-9001-B','SHARD_2','FAILED','PENDING');COMMIT;SELECT *FROM sh26_distributed_opsORDER BY shard_name;
This is a state-machine teaching aid, not Oracle's actual 2PC implementation. It makes the application design question explicit: what should the business API return while global outcome is unresolved?
8. Deliberately wrong: blindly retry a distributed payment/update
If the client loses its network response after commit, it may not know whether the database committed. Blindly repeating a non-idempotent request can double-charge/create duplicate work. The repair is a durable idempotency/request key whose business result can be safely returned again.
CREATE TABLE sh26_request_dedup ( request_id VARCHAR2(64) PRIMARY KEY, tenant_id NUMBER NOT NULL, operation_name VARCHAR2(40) NOT NULL, outcome_code VARCHAR2(20), created_at TIMESTAMP DEFAULT SYSTIMESTAMP NOT NULL);INSERT INTO sh26_request_dedup( request_id,tenant_id,operation_name,outcome_code)VALUES('REQ-9001',1,'CROSS_TENANT_ADJUSTMENT','ACCEPTED');COMMIT;-- A retry first checks the same durable request_id.SELECT outcome_codeFROM sh26_request_dedupWHERE request_id='REQ-9001';
9. Design APIs around locality
| API | Preferred scope | Why |
|---|---|---|
| Get/update work order | Single shard by tenant | Low latency and one local transaction. |
| Create technician + work order | Same shard/table family | Keep referential/business transaction local. |
| Global revenue report | Explicit multi-shard | Fan-out/aggregation is inherent and can have separate SLO. |
| Cross-tenant administrative adjustment | Distributed/saga-like workflow only when required | Needs 2PC or explicit compensating/idempotent semantics. |
10. Failure domains change application error handling
One shard can be down while other shard-key requests succeed. A global query may fail or degrade because one participant is unavailable. A local retry policy should not accidentally turn a one-shard outage into a retry storm against every shard. Expose shard/key/request identity in telemetry and use bounded backoff/circuit-breaking at the application/service layer.
11. Deliberately wrong: assume all cross-shard SQL shapes behave like local Oracle SQL
Distributed query processing has explicit supported-shape
restrictions. A query using a locally valid feature such as
CONNECT BY can be unsupported as a multi-shard
query. The repair is to check distributed-query support, push
shard-local work where supported, or redesign the reporting
pipeline/API.
12. Cleanup
DROP TABLE sh26_request_dedup PURGE;DROP TABLE sh26_distributed_ops PURGE;DROP TABLE sh26_xshard_orders PURGE;
13. Production judgment
Cross-shard work is not “bad,” but it belongs to a different SLO/correctness class. Favor key-routed local transactions for interactive OLTP; make global reports explicit; use distributed transactions only when atomic cross-shard business semantics justify prepare/commit/in-doubt complexity. Add idempotency keys and telemetry for ambiguous client outcomes.
26ai supports parallel cross-shard DML, but it does not remove
2PC/failure-domain costs. Genuine in-doubt state is inspected
through distributed-transaction views and resolved by Oracle
recovery/DBA procedures. No local lab changes
ALLOW_LEGACY_RECO_PROTOCOL. Lesson 4 now treats the
distributed topology itself as mutable: adding shards, moving
chunks and enforcing regional placement are operational
migrations.
Check your understanding
- What additional work does a multi-shard query coordinator perform?
- Why can a distributed transaction become in doubt?
- Does 26ai parallel cross-shard DML remove two-phase commit?
- Why are idempotency keys important after a network timeout?
- What is the safest latency class for interactive sharded OLTP?
Review the answers
It fans statements to shards, receives partial results and can perform final merge/aggregation/SCN coordination.
A failure can interrupt prepare/commit communication after participants have durable prepared state.
No. Parallelism improves execution scheduling; atomic commit still requires distributed coordination.
The client may not know whether the first request committed; the same request ID makes a retry return/reconcile the existing outcome instead of duplicating work.
A single-shard key-routed transaction whenever the business operation can be modeled that way.
Authoritative references
- Query Processing for Multi-Shard Queries — fan-out/coordinator processing
- Supported Query Constructs — distributed query restrictions
- Supported DMLs and Examples — cross-shard DML/2PC
- Distributed Transactions Concepts — 2PC/in-doubt mechanism
- Changes in 26ai — parallel cross-shard DML