A distributed transaction is a unit of work whose operations span two or more independent data stores, services or network nodes, yet must complete with all-or-nothing atomicity across every participant. Coordinating such a transaction requires protocols that agree on a single outcome despite partial failures, network partitions and concurrent activity at each site. Classical coordination uses atomic commit protocols, while modern systems often relax strict atomicity for availability using compensating workflows.

Overview

  • Distributed transactions extend the familiar database guarantee of atomicity, consistency, isolation and durability to settings where state is partitioned across machines that can fail independently. Because no single node has authoritative knowledge of all participants, the system relies on a coordinator and an agreement protocol to drive every participant to the same commit or abort decision.
  • It is modelled as a subclass of Transaction Processing within the distributed-systems domain.
  • The fundamental tension in distributed transactions is captured by the trade-off between strong atomicity and availability under partition. Two-phase commit delivers atomicity but blocks if the coordinator fails at the wrong moment, while consensus-backed commit and saga-based compensation relax timing or isolation to keep services responsive during failures.
  • Modern microservice and ledger systems frequently prefer eventual, compensating designs over global locks because cross-service two-phase commit couples availability across teams and scales poorly. Sagas trade the clean abstraction of a single atomic unit for a sequence of locally durable steps with explicit undo logic.

Mechanisms

  • Atomic commit: two-phase and three-phase commit protocols vote and then converge on a single global decision.
  • Coordinator failure: blocking in two-phase commit motivates consensus-backed coordination and recovery logs.
  • Saga compensation: long-running business transactions trade strict atomicity for a sequence of locally committed steps with compensating undo actions.
  • Isolation levels: concurrency control and global ordering manage interleaving across participants.

Applications

  • Cross-shard updates in distributed and partitioned databases.
  • Multi-service workflows in microservice architectures using sagas.
  • Atomic settlement and cross-ledger exchange in financial and blockchain systems.

Considerations

  • Isolation across participants is hard; without care, intermediate states of a saga become visible and require careful semantic compensation.
  • Coordinator and participant recovery depend on durable logs so that an interrupted transaction can be resolved deterministically on restart.
  • Idempotency of participant operations is a practical prerequisite for safe retries in any distributed commit protocol.

Provenance