Skip to content

Archive

Distributed Systems

108 articles
Software Engineering 15 Sep 2026 6 min read

Fencing Tokens Reject Stale Lease Holders

A distributed lease can expire while its holder is paused. The holder may later resume with local state that still says it owns the lease, even though another client has acquired a newer lease. If the protected resource accepts commands solely because a client once acquired ownership, both clients can act during the same logical ownership interval. Fencing tokens move part of the ownership check to the resource receiving the mutation. Each successful lease acquisition receives a token greater than every token issued before it. The protected resource records the greatest token it has accepted and rejects operations carrying an older value.

Software Engineering 14 Sep 2026 7 min read

Transactional Outboxes Move Atomicity Into the Database

A service updates an order row and emits an event about that update. If the database commit succeeds but the broker publish fails, durable state says one thing while downstream consumers receive no corresponding message. Reversing the order only moves the gap: a successful publish followed by a failed database transaction exposes an event for state that never committed. The difficulty is not message syntax or retry configuration. It is atomicity across two systems that do not share a transaction. A transactional outbox changes the boundary. The application writes its domain state and a message record into the same database transaction, then a separate relay publishes committed outbox records to the broker.

Software Engineering 14 Sep 2026 7 min read

Poison Messages Turn Retries Into Queue Retention

A queue consumer receives a message, rejects it, and receives the same message again. That cycle is useful when the rejection came from a transient condition. It is structurally different when the payload can never be processed by the current consumer. The broker can keep honoring redelivery semantics while the application makes no forward progress on that message. Such a message is commonly called a poison message. The important property is not that it contains malformed bytes. A syntactically valid message can be permanently unprocessable because its schema is unsupported, a required invariant is violated, referenced data can never exist, or application logic deterministically rejects its state.

Software Engineering 14 Sep 2026 9 min read

Deadlines Shrink Across Service Boundaries

A service receives a request with 480 milliseconds remaining before its deadline. It spends 90 milliseconds reading state, then calls another service with a fixed 500-millisecond timeout. The downstream call can now outlive the request that caused it. Nothing about either timeout is internally inconsistent; the inconsistency appears at the boundary between them. Timeouts are often configured as local limits: a database query gets one value, an HTTP client another, a queue operation a third. A deadline represents a different constraint. It gives an operation an end point, so every later stage can compare its own work against the same finite lifetime.

Software Engineering 13 Sep 2026 7 min read

Idempotent Consumers and Durable Duplicate Detection

Idempotent Consumers and Durable Duplicate Detection A consumer commits a database transaction, then loses its connection before acknowledging the message. The broker has no evidence that processing finished, so a delivery protocol that permits redelivery can present the same message again. The second delivery is not evidence that the first transaction failed. From the consumer’s perspective, the important fact is more precise: message delivery and application commit have separate completion points. If the broker cannot atomically participate in the application’s state transition, an acknowledgement can be lost after the application effect is already durable.

Software Engineering 13 Sep 2026 6 min read

Fencing Tokens Make Expired Leases Observable

A process can hold a distributed lease, pause long enough for that lease to expire, then resume with local state that still says it owns the resource. Another process may already have acquired a newer lease during the pause. At that point, mutual exclusion in the lock service is not enough: two processes can each act as if they have authority, even though only one lease is current. This is a boundary problem between coordination and the resource being protected. A lease service can decide which holder is current according to its own state. It cannot retroactively erase instructions already held by an old process, nor can it stop that process from sending a request after a long pause.

Software Engineering 13 Sep 2026 8 min read

Circuit Breakers Bound Failed Call Admission

A remote call can fail in a few milliseconds or consume its entire timeout budget before returning an error. If callers keep issuing equivalent requests while the dependency remains unable to serve them, each attempt spends resources on an outcome that recent evidence already suggests is unavailable. Retries can increase that pressure because one logical operation may create several physical calls. A circuit breaker changes call admission rather than the remote protocol. It records recent outcomes, moves between explicit states, and can reject new calls locally for a bounded period. After that period, it permits limited probes to test whether normal traffic can resume.

Software Engineering 12 Sep 2026 8 min read

Vector Clocks and the Shape of Concurrent State

A single version number can state that one value came after another only when every update participates in the same ordered sequence. Replicated state breaks that assumption as soon as independent writers can accept changes without first agreeing on one global next version. Two replicas can each move forward from the same ancestor. Calling one state version 8 and the other version 9 creates an order, but that order may describe the numbering scheme rather than the causal relation between the writes. Vector clocks represent a different fact: which update history a state has observed.

Software Engineering 12 Sep 2026 7 min read

Transactional Outbox: Moving Publication Across the Commit Boundary

Transactional Outbox: Moving Publication Across the Commit Boundary A service that changes database state and publishes a message has two distinct side effects. A local transaction can make the database change atomic, and a broker can accept the message durably, but those facts do not make the pair atomic. The awkward interval sits between them. If the database commit succeeds and publication does not, durable state exists without the corresponding message. Reversing the order only reverses the exposure: a message can become visible before the database commit succeeds.

Software Engineering 12 Sep 2026 9 min read

Transactional Outbox at the Database-Broker Boundary

Transactional Outbox at the Database-Broker Boundary A database commit and a broker publish are two separate state transitions. An application can complete either one first, but unless both systems participate in a common transaction protocol, there is an interval in which one side has changed and the other has not. That interval is the central problem behind the transactional outbox pattern. The pattern does not make a database and broker commit atomically. Instead, it moves the durable publication decision into the same database transaction as the application state change. A separate publisher later converts that recorded intent into a broker message.

Software Engineering 12 Sep 2026 12 min read

Saga Transactions: Coordinate Multi-Service Changes with Compensation

Saga Transactions: Coordinate Multi-Service Changes with Compensation A business operation can cross several services even when no single database transaction spans them all. An order flow might reserve inventory, authorize payment, create a shipment, and confirm the order. Each service owns its data and commits independently. That independence creates a difficult failure case. Inventory can be reserved successfully, then payment authorization can fail. A database rollback in the order service cannot undo a committed reservation in the inventory service.

Software Engineering 12 Sep 2026 8 min read

Retry Amplification and the Role of Jitter

Retry Amplification and the Role of Jitter A failed request can create more traffic than a successful one. If a caller immediately repeats an operation after a transient error, the original unit of demand becomes two attempts. Add another retrying layer above that caller, and a single logical request can fan out into several physical attempts before any component has recovered. Retries are often described as a way to tolerate temporary faults. That description is incomplete because retry behavior also changes load. The mechanism sits inside a feedback loop: failure triggers another attempt, another attempt consumes capacity, and consumed capacity can affect the conditions that produced the failure.

Software Engineering 12 Sep 2026 10 min read

Request Coalescing at Hot Cache Misses

Request Coalescing at Hot Cache Misses A cache entry expires at 22:00:00.000. Ten milliseconds later, two hundred requests ask for the same key. A cache-aside implementation sees two hundred misses. If every caller independently reads the backing service, a single expiration event becomes two hundred concurrent backend operations. Nothing is wrong with the cache lookup itself. The amplification comes from treating identical in-flight work as unrelated. Request coalescing changes that boundary. Callers that need the same absent key share one active load, while requests for other keys continue independently. The mechanism is small, but its semantics reach beyond a mutex: it defines which operations may share a result, how failures fan out, what cancellation means, and when another load may begin.

Software Engineering 12 Sep 2026 11 min read

Idempotent Consumers: Handle Duplicate Messages Safely

Idempotent Consumers: Handle Duplicate Messages Safely A message broker can deliver the same message more than once. A worker may finish its database update and crash before acknowledging the message. The broker sees no acknowledgement, so it sends the message again. From the broker’s perspective, redelivery is the safe choice. From the application’s perspective, the second delivery can repeat a business effect. That gap matters whenever an effect must happen once per logical message. Charging an account twice, granting stock twice, incrementing a counter twice, or sending the same fulfillment request twice can turn a routine retry into corrupted state.

Software Engineering 12 Sep 2026 9 min read

Fencing Tokens: Block Stale Lease Holders

Fencing Tokens: Block Stale Lease Holders A distributed lease can grant one process temporary permission to act, but expiration alone cannot stop that process from acting after its lease has ended. A long pause, network delay, overloaded runtime, or suspended virtual machine can leave an old holder unaware that another process has already acquired the lease. This creates a subtle safety gap. Two processes can both believe they are entitled to modify the same resource, even when the lease service itself grants ownership correctly.

Software Engineering 12 Sep 2026 8 min read

Fencing Tokens for Expiring Distributed Leases

A lease can expire while its holder is still running. That single property separates a distributed lease from an ordinary in-process mutex. The coordinator may grant ownership to another client after a deadline, yet the former holder can resume after a long pause and continue issuing operations based on authority it no longer has. The coordinator has done its job: it stopped treating the old client as the current holder. The shared resource has a different problem. Unless operations carry evidence of ownership order, the resource may have no basis for distinguishing a current holder from a stale one.

Software Engineering 12 Sep 2026 7 min read

Deadline Propagation as a Request Boundary

Deadline Propagation as a Request Boundary A service can return after its caller has stopped waiting. The computation may still consume a connection, hold a concurrency slot, execute a database query, or start another remote call. A local timeout limits how long one caller waits; it does not, by itself, bound the lifetime of work already sent deeper into the system. An end-to-end deadline changes that boundary. Instead of giving each operation an independent duration, the request carries a point in time after which its result is no longer useful to the initiating operation. Each component can derive its remaining budget from that same boundary.

Software Engineering 12 Sep 2026 9 min read

Consistent Hashing: Limit Key Movement as Nodes Change

Distributed systems often need a deterministic answer to a simple question: given a key, which node should own it? A cache cluster may route each object key to one server. A storage service may assign each partition to a shard. A worker pool may send all events for the same account to the same processor. The routing rule must be stable enough that clients agree, yet flexible enough to handle nodes joining and leaving.

Software Engineering 12 Sep 2026 7 min read

Compensation Is Not Rollback Across Service Boundaries

Compensation Is Not Rollback Across Service Boundaries A local database rollback can erase uncommitted writes before other transactions are allowed to depend on them. A compensating operation has a different shape. It runs after an earlier operation has committed, often after that result has become visible to other components. That distinction changes the consistency model. Compensation does not restore a distributed system to a state in which the original action never occurred. It adds another state transition whose domain meaning offsets some consequence of the first one.

Software Engineering 12 Sep 2026 9 min read

Circuit Breakers as Admission Control for Failing Dependencies

Circuit Breakers as Admission Control for Failing Dependencies A remote call that has little chance of succeeding still consumes something: a connection slot, a worker, a deadline budget, memory for request state, or capacity in the dependency itself. When repeated failures indicate that a downstream service is currently unable to serve useful work, continuing to admit every call can preserve the very pressure that callers need to escape. A circuit breaker changes that admission decision. Instead of treating each call as independent, it retains a small amount of state about recent outcomes. That state can temporarily reject new calls before network I/O begins, then permit controlled probes after a recovery interval.

Software Engineering 12 Sep 2026 11 min read

Bulkhead Isolation: Contain Failures with Separate Capacity Pools

Bulkhead Isolation: Contain Failures with Separate Capacity Pools A service can have enough total capacity and still become unavailable because one dependency consumes all of it. Imagine an API that calls a payment service and a recommendation service. Both outbound calls use the same worker pool. Recommendations become slow during a traffic spike. Their requests occupy every worker while waiting for responses. Payment requests now have no worker available, even though the payment service itself is healthy.

Software Engineering 12 Sep 2026 10 min read

Backpressure: Match Producer Speed to Consumer Capacity

Backpressure: Match Producer Speed to Consumer Capacity A pipeline is stable only when work enters at a rate its downstream stages can sustain. That sounds obvious, yet many systems let producers run at full speed until a queue fills, memory grows, latency explodes, or a downstream service starts rejecting requests. The visible failure appears late. The actual mismatch began earlier: one stage could create work faster than the next stage could finish it.

Software Engineering 11 Sep 2026 9 min read

Fencing Tokens: Stop Stale Lock Holders from Writing

Fencing Tokens: Stop Stale Lock Holders from Writing A distributed lock can tell a client that it owns a resource for a limited period. That does not guarantee the client stops acting when the period ends. A process can pause for garbage collection, lose network access, become descheduled, or stall on an overloaded machine. During that pause, its lease can expire and another client can acquire the same lock. When the first process resumes, it may still believe it is entitled to write.

Software Engineering 11 Sep 2026 8 min read

Consumer-Driven Contract Tests for Service Compatibility

Consumer-Driven Contract Tests for Service Compatibility Two services can pass their own test suites and still fail when deployed together. A provider may rename a field, narrow an accepted value, change a status code, or remove an endpoint. Its internal tests can remain green because those tests describe the provider’s own view of correct behavior. A consumer can still depend on the old interaction. Consumer-driven contract testing turns selected consumer expectations into executable contracts. The consumer records the interactions it requires. The provider then verifies those contracts against its implementation.