Skip to content

Archive

Distributed Systems

108 articles
Artificial Intelligence 24 Sep 2026 5 min read

Expert Parallelism Turns MoE Routing into All-to-All Communication

A sparsely activated Mixture-of-Experts layer can keep many expert parameters while sending each token to only a small subset of them. Once those experts are partitioned across devices, sparse activation does not imply local execution. A token assigned to a remote expert has to move to the device that owns that expert, and its expert output has to return to the device that continues the model computation. This movement is the defining systems cost of expert parallelism. The router makes a logical selection, but a distributed runtime must turn that selection into communication, local expert batches, expert computation, and a reverse exchange.

Software Engineering 22 Sep 2026 6 min read

Transactional Outbox Closes the Database-to-Broker Commit Gap

Transactional Outbox Closes the Database-to-Broker Commit Gap A service often needs one operation to change database state and emit a message. An order may become confirmed while an OrderConfirmed event is sent to a broker. Those actions touch separate systems, so two ordinary writes cannot form one atomic commit unless both systems participate in a distributed transaction. The dangerous part is the interval between the writes. Commit the database first and the process can fail before publishing. Publish first and the database commit can fail afterward. Reversing the order moves the failure window; it does not remove it.

Software Engineering 22 Sep 2026 6 min read

Tombstones Prevent Deleted Data from Reappearing

Tombstones Prevent Deleted Data from Reappearing Deletion is not merely the absence of a value in a replicated store. Absence carries no information about whether a key was deliberately removed or whether a replica has simply never received it. When replicas can be temporarily disconnected, that distinction determines whether synchronization preserves a deletion or accidentally restores old data. A tombstone records the deletion as versioned state. Replicas can compare that marker with older values and keep the deletion when they reconcile. The marker can eventually be reclaimed, but only after the system has a defensible boundary beyond which an older value cannot return.

Software Engineering 22 Sep 2026 7 min read

Saga Compensation Is Not Transaction Rollback

Saga Compensation Is Not Transaction Rollback A multi-service operation can cross inventory, payments, shipping, and other independently committed systems. Once one service commits its step, a later failure cannot make that earlier commit disappear through an ordinary database rollback. A saga handles this boundary by pairing forward actions with explicit recovery actions. If a later step fails, the coordinator invokes compensations for earlier completed steps where the business process permits them.

Software Engineering 22 Sep 2026 6 min read

Rendezvous Hashing Limits Key Movement During Membership Changes

Rendezvous Hashing Limits Key Movement During Membership Changes A partitioning rule has two jobs that can pull in different directions. It should spread keys across available nodes, and it should avoid moving most keys when that node set changes. A simple modulo rule handles the first job well for a stable cluster but performs poorly at the second. Rendezvous hashing, also called highest-random-weight hashing, assigns every key a deterministic score for every eligible node. The node with the highest score owns the key. Adding or removing a node changes only the comparisons involving that member, so keys with unaffected winners keep their placement.

Software Engineering 22 Sep 2026 5 min read

Load Shedding Protects Services When Capacity Runs Out

Load Shedding Protects Services When Capacity Runs Out A service can receive more work than it can complete. The first visible symptom is often not an immediate error but a queue that grows while workers remain fully occupied. Requests spend longer waiting, deadlines expire, clients retry, and the extra retry traffic can deepen the overload. Load shedding places an explicit rejection point before that spiral consumes every available resource. The service admits work that fits its operating capacity and fails excess work quickly enough to preserve useful throughput for requests that can still complete.

Software Engineering 22 Sep 2026 6 min read

Idempotency Keys Turn Retries into One Logical Operation

Idempotency Keys Turn Retries into One Logical Operation A client can lose the response to a successful request. The server may commit a charge, create an order, or enqueue a job, then the connection can fail before the response reaches the caller. From the client’s view, success and failure are now ambiguous. Retrying is necessary for availability, but an ordinary retry can repeat the side effect. An idempotency key gives the client a way to say that several HTTP attempts represent one logical operation.

Software Engineering 22 Sep 2026 8 min read

Hedged Requests Trade Duplicate Work for Lower Tail Latency

Hedged Requests Trade Duplicate Work for Lower Tail Latency Most requests may finish quickly while a small fraction take much longer. A busy worker, a transient network queue, a cold cache entry, garbage collection, storage contention, or another local disturbance can stretch one attempt far beyond the median. At scale, those slow outliers become visible in p95, p99, and higher-percentile latency even when average service time looks healthy. A hedged request sends an additional attempt after the original has been outstanding for a chosen delay. Both attempts represent the same logical operation. The caller accepts the first valid result and cancels or ignores the remaining attempt.

Software Engineering 22 Sep 2026 6 min read

Fencing Tokens Stop Stale Lease Holders

Fencing Tokens Stop Stale Lease Holders A distributed lease gives one worker temporary permission to act as an owner. The lease eventually expires so another worker can take over after a crash or network failure. That solves availability, but expiry alone does not guarantee that the old worker has stopped. A process can pause long enough for its lease to expire, then resume with stale local state. A long garbage-collection pause, scheduler stall, suspended virtual machine, or delayed network path can create this condition. If the old worker writes after a replacement has taken ownership, two workers can affect the same resource even though the lease service never considered both leases valid at the same instant.

Software Engineering 22 Sep 2026 7 min read

Fencing Tokens Block Stale Lock Holders

Fencing Tokens Block Stale Lock Holders A distributed lock is often used to keep two workers from changing the same resource at once. The difficult case begins when lock ownership depends on a lease. A client can acquire the lease, pause long enough for it to expire, then resume after another client has acquired a new lease. From the old client’s point of view, execution simply continued. From the coordination service’s point of view, ownership already moved. If the protected storage system accepts both clients’ writes, the old holder can overwrite work performed by the current holder.

Software Engineering 22 Sep 2026 6 min read

Dead-Letter Queues Isolate Poison Messages Without Blocking Progress

Dead-Letter Queues Isolate Poison Messages Without Blocking Progress A message consumer usually treats failure as temporary at first. A database may be unavailable, a remote service may time out, or a worker may restart between receiving and acknowledging a message. Retrying is appropriate when another attempt has a reasonable chance of succeeding. Some messages fail for a different reason. Their payload is malformed, a referenced entity can never satisfy a required condition, or the consumer has a deterministic defect triggered by that input. Repeated delivery then consumes capacity without moving the message toward completion. A dead-letter queue gives that failure a separate destination after the normal retry policy is exhausted.

Software Engineering 22 Sep 2026 7 min read

Consistent Hashing Limits Key Movement During Topology Changes

Consistent Hashing Limits Key Movement During Topology Changes A distributed cache or partitioned service needs a rule that maps each key to a node. A simple rule such as hash(key) % N is attractive while the node count stays fixed. The trouble appears when N changes. Moving from four nodes to five changes the divisor for every key. Most remainders change, so a routine capacity adjustment can remap a large share of the dataset at once. For a cache, that can trigger a wave of misses. For stateful storage, it can create a large migration job.

Software Engineering 22 Sep 2026 7 min read

Circuit Breakers Stop Repeated Calls to Failing Dependencies

Circuit Breakers Stop Repeated Calls to Failing Dependencies A dependency that is already failing can consume more caller capacity than a healthy one. Requests wait for timeouts, retries add traffic, connection pools remain occupied, and worker slots stay tied to work that has little chance of completing. A circuit breaker places a stateful gate in front of that dependency so the caller can stop issuing calls after failure evidence reaches a configured limit.

Software Engineering 22 Sep 2026 6 min read

Bulkheads Keep One Saturated Dependency from Consuming Every Worker

Bulkheads Keep One Saturated Dependency from Consuming Every Worker A service can have healthy CPU, available memory, and responsive internal code while still becoming unavailable. One downstream dependency is enough to consume the service’s entire concurrency budget if calls to it become slow and every request is allowed to wait. The failure is not limited to the slow dependency. Shared worker pools, connection pools, semaphores, queues, and request slots turn local saturation into a service-wide outage. Bulkhead isolation limits that blast radius by reserving separate capacity for distinct workloads or dependencies.

Software Engineering 22 Sep 2026 7 min read

Bounded Queues Turn Overload into an Explicit Admission Decision

Bounded Queues Turn Overload into an Explicit Admission Decision A queue absorbs short differences between arrival rate and service rate. That buffer is useful when a burst ends before workers fall far behind. The same mechanism becomes dangerous when arrivals remain faster than completions: every accepted item adds waiting time and consumes some combination of memory, descriptors, references, or durable storage. A bounded queue places a finite limit on that waiting population. Once the limit is reached, the system must make an admission decision instead of silently extending the backlog. Depending on the interface, that decision may block a producer, reject new work, shed selected work, or redirect it to another capacity domain.

Software Engineering 21 Sep 2026 7 min read

Version Vectors Separate Causality from Concurrency

Version Vectors Separate Causality from Concurrency Replicated data can receive writes at several nodes while communication between those nodes is delayed. When two versions meet later, a store has to decide whether one descends from the other or whether both were created independently. A wall-clock timestamp gives a total-looking order, but clock order is not causal order. Two replicas can write during a partition, and whichever timestamp happens to be larger does not make that write a descendant of the other.

Software Engineering 21 Sep 2026 4 min read

Version Vectors Distinguish Concurrent Updates from Causal Successors

Version Vectors Distinguish Concurrent Updates from Causal Successors Replicated data can receive writes at different nodes while communication between those nodes is delayed. When versions later meet, a scalar revision number can say that two values differ, but it cannot always say whether one descends from the other or both were produced independently. A version vector records progress per replica. Comparing those counters provides a partial order: one version can dominate another, the vectors can be equal, or neither can dominate. The last case identifies concurrent histories that require an explicit reconciliation rule.

Software Engineering 21 Sep 2026 8 min read

Transactional Outbox Closes the Dual-Write Gap

Transactional Outbox Closes the Dual-Write Gap A service often needs one request to change database state and publish an event. The two operations may look adjacent in application code, but they cross different durability boundaries. A database commit can succeed while a broker publish fails, or the publish can succeed before the database transaction rolls back. That split creates a dual-write problem. No ordering of two independent writes can make them atomic by itself.

Software Engineering 21 Sep 2026 5 min read

Transactional Outbox Closes the Database-Broker Commit Gap

Transactional Outbox Closes the Database-Broker Commit Gap A service often needs one request to change database state and publish a message. Those actions may look adjacent in application code, but they cross two independent commit boundaries. If the database and broker do not share a transaction protocol, no ordering of two ordinary writes can make them atomic. Consider an order service that stores an accepted order and emits OrderCreated. Publishing after the database commit leaves a crash window before the broker call. Publishing first creates the opposite window: consumers can receive an event for state that later fails to commit.

Software Engineering 21 Sep 2026 6 min read

Token Buckets Preserve Burst Capacity Without Removing Rate Bounds

Token Buckets Preserve Burst Capacity Without Removing Rate Bounds A fixed requests-per-second ceiling treats a brief spike and a sustained flood as the same event. That can be too rigid for services whose callers naturally arrive in clusters. A token bucket separates two constraints: the long-run admission rate and the amount of burst traffic the service is willing to absorb. The model has two parameters. The bucket capacity B is the maximum number of tokens that can accumulate. The refill rate r adds tokens per unit of time, up to B. A request consumes tokens according to its configured cost. If enough tokens are present, the request proceeds; otherwise it is rejected, delayed, or handled by another explicit policy.

Software Engineering 21 Sep 2026 7 min read

Load Shedding Protects Useful Work Under Saturation

Load Shedding Protects Useful Work Under Saturation A service has a finite amount of work it can complete per unit of time. When offered load rises past that capacity, accepting every request does not create more capacity. It creates more waiting, consumes memory and connection slots, extends deadlines, and can leave expensive work running after callers have already given up. Load shedding makes admission explicit. Work that the service cannot process within its operating envelope is rejected early so admitted work retains a realistic chance of completing.

Software Engineering 21 Sep 2026 7 min read

Idempotency Keys Make Retried Mutations Safe

Idempotency Keys Make Retried Mutations Safe A client can lose the result of a successful mutation without losing the mutation itself. The server may commit a payment, reservation, or job submission and then drop the connection before the response reaches the caller. From the client’s perspective, a timeout leaves two plausible states: the operation failed before commit, or it committed and only the response was lost. Blindly retrying a non-idempotent mutation can apply the effect twice. Refusing every retry leaves the caller unable to recover from an ambiguous outcome. An idempotency key gives both sides a stable identity for one logical operation, so a repeated attempt can reuse the result of the first accepted attempt instead of creating another effect.

Software Engineering 21 Sep 2026 6 min read

Idempotency Keys Bound Duplicate Mutations Across Retries

Idempotency Keys Bound Duplicate Mutations Across Retries A client can lose the response to a successful mutation. The server may commit a payment, reservation, or job submission and then lose the connection before the response reaches the caller. From the client side, timeout does not reveal whether the mutation failed before commit or succeeded before the response disappeared. A retry is necessary for availability, but a blind retry can repeat the side effect. An idempotency key gives both attempts a stable identity so the server can treat them as one logical operation.

Software Engineering 21 Sep 2026 7 min read

Hedged Requests Cut Tail Latency with Controlled Duplication

Hedged Requests Cut Tail Latency with Controlled Duplication A service can have a healthy median latency while a small fraction of requests take far longer. Queueing, garbage collection, storage stalls, packet loss, noisy neighbors, or uneven replica load can all stretch the slow end of the distribution. For a request that fans out to several dependencies, one slow branch can dominate the entire response. Hedged requests reduce that exposure by starting a second copy after a short delay. The copies target independent execution paths when possible, and the first valid response wins.