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.

Advanced130–150 minutesCross-shard fan-out + 2PC failure simulation26ai parallel cross-shard DML awarenessIn-doubt transactions are operational stateLast reviewed: August 2026

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.

01

Explain query coordinator fan-out/partial aggregation and why single-shard SQL is a different latency class.

02

Explain global read consistency/SCN coordination for multi-shard reads and current consistency modes at architecture level.

03

Explain 26ai cross-shard DML and two-phase commit, including in-doubt transaction recovery.

04

Design locality-preserving, idempotent application APIs with explicit retry/request IDs.

05

Simulate partial shard failure and distributed transaction state on one Free database without pretending it is actual 2PC.

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. 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.

sql · real conceptual cross-shard query
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.

sql · real coordinator-side DML shape
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.

sql · distributed transaction evidence on Oracle databases
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.

sql · read-only preflight
SELECT name,value,issys_modifiable,ispdb_modifiableFROM v$parameterWHERE name='allow_legacy_reco_protocol';

6. Free simulation: partial aggregation across logical shards

sql · setup
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;
sql · simulate shard-local partial aggregation
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

sql · coordinator state machine table
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.

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

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

  1. What additional work does a multi-shard query coordinator perform?
  2. Why can a distributed transaction become in doubt?
  3. Does 26ai parallel cross-shard DML remove two-phase commit?
  4. Why are idempotency keys important after a network timeout?
  5. 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

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.