A distributed system rarely fails with a clean error message. The more dangerous failure is uncertainty.
A customer clicks Pay. The request reaches a payment provider, but the response times out. Did the payment fail? Did it succeed but the response disappear? Is it still running? A simple retry can charge the customer twice. An asynchronous event can arrive before the order service has recorded the payment. A queue can grow while every dashboard remains reassuringly green.
This is where distributed systems engineering begins. Not with a diagram of services and arrows, but with a decision about what the system is allowed to do when it cannot know the full truth.
This guide follows that order flow through the failure modes that matter in production, then explains the engineering controls that make recovery safe.
Start with one rule: the network is not a function call
Inside one process, a function either returns a result or throws an error. Across a network, a caller may see a timeout while the receiver completes the work. Messages can be delayed, duplicated, reordered or lost. Components can disagree temporarily about the same business event.
That does not mean every application needs Kafka, multi-region writes or a distributed consensus system. A well-designed modular monolith with one primary database is often the safer choice. Distribution earns its complexity only when independent scaling, availability, geography or ownership genuinely require it.
The scenario: one payment, many possible states
Imagine an ecommerce order service. It accepts an order, asks a payment provider to authorise the amount, records the result, publishes an event, reserves inventory and emails the customer.
At least five systems now participate in what a customer experiences as one action. The engineering question is no longer just “did the API return 200?” It is “which state is authoritative, which work can be repeated, and how do we repair disagreement?”
1. Timeout ambiguity
A timeout only tells the caller that it did not receive an answer within the time budget. It does not prove that the receiver did nothing.
If the order service retries an unknown payment request with a new identifier, the provider may process both calls. If the service marks the order failed without checking, it may reject a legitimate purchase.
Solution: durable intent and idempotency keys
Record the customer’s intended operation before making the external call. Assign an idempotency key that represents the business action, not the individual HTTP attempt. Send the same key whenever that action is retried. The payment boundary should return the original result for a repeated key rather than make a second charge.
Then model the order as an explicit state machine, for example payment_pending, payment_authorised, payment_declined and payment_review. Unknown is a real state. Treating it as failure creates expensive customer and operational errors.
2. Duplicate messages and duplicate work
At-least-once delivery is common because it is practical. A broker may redeliver a message after a worker crashes, or a publisher may retry after losing the acknowledgement. This is not a bug in the queue. It is a condition the consumer must be designed to survive.
A duplicate payment_authorised event must not reserve stock twice, send two receipts or create two fulfilment jobs.
Solution: make consumers idempotent
Give every event a stable event ID and retain a processed-event record at the consumer boundary. When a worker receives the same ID again, it returns safely without repeating the business effect. Where practical, combine the state change and the processed-event marker in one database transaction.
Do not rely on “exactly once” as a business guarantee. The useful guarantee is that repeating the operation does not change the intended outcome after the first successful application.
3. Retry storms
Retries help only when the failure is temporary and the dependency has room to recover. During an outage, thousands of clients retrying together can turn a slow service into an unavailable one. The recovery path becomes the source of the incident.
Solution: time budgets, bounded retries and jitter
Set a clear deadline for the complete user operation, then allocate smaller time budgets to dependencies. Retry only errors that are plausibly transient. Keep retry attempts small. Use exponential backoff with random jitter so callers do not return in lockstep.
When an error rate crosses a threshold, a circuit breaker should stop sending work to the unhealthy dependency for a short period. Return a useful pending state, use a queue for work that does not need a synchronous answer, or fall back to a reduced but truthful experience. The right fallback is not one that hides failure. It is one that preserves the customer’s intent without pretending the work is complete.
4. Queues that become invisible backlog
Moving work to a queue removes it from the request path. It does not remove the work. If order events arrive faster than workers can process them, the system has converted latency into waiting time.
A queue depth alone can mislead because traffic changes. A short queue of old jobs can be worse than a deep queue of new ones.
Solution: backpressure and queue-age ownership
Measure the age of the oldest message, processing rate, failure rate and retries by error class. Place limits on producer rate or worker concurrency when a downstream dependency is saturated. Send repeatedly failing messages to a dead-letter queue with enough context for a human or repair job to act.
Every queued workflow needs an owner and a recovery rule. “It is asynchronous” is not an operational state.
5. Events that arrive out of order
Network timing does not preserve business meaning. An inventory release can arrive after a later reservation. A customer can receive a status update from a slower replica after seeing a newer one. A worker can process version three before version two.
Solution: versioned state transitions
Attach a version, sequence number or business timestamp to state changes where ordering matters. Accept a transition only when it is valid from the current state. For example, an order should not move from shipped back to payment_authorised because a delayed event arrived.
Where a total order is unnecessary, define a conflict rule deliberately. A product catalogue can tolerate eventual agreement. A payment balance or access decision usually cannot. Consistency is not a feature to maximise everywhere. It is a business decision about where disagreement is unacceptable.
6. Competing writers and split brain
The most serious state is usually not owned by every service. It has one authority. If two components both believe they can assign the same inventory, issue the same refund or lead the same control plane, the system can corrupt data while appearing available.
Solution: one authority, fencing and quorum when required
Keep a single authoritative writer for critical state whenever the business permits it. For leader-based coordination, use a lease or fencing token so an old leader cannot continue writing after losing authority. When several replicas must agree on critical state, use a proven coordination system with quorum rules rather than inventing consensus inside application code.
Google’s SRE guidance describes split-brain as a critical shared-state problem: network partitions can lead each side to elect a leader and accept changes unless the system requires sufficient healthy replicas and connectivity to make progress safely. Read the SRE discussion of consensus and split brain.
7. Observability without causality
Most incidents contain logs. The hard part is showing how one customer action moved through an API, queue, worker, database and third-party dependency. Without that thread, teams guess. Guesses produce broad retries, unnecessary rollbacks and prolonged incidents.
Solution: make one request traceable end to end
Create a correlation ID at the first boundary and carry it through synchronous calls, messages and background jobs. Emit structured logs with the operation ID, business entity ID, state transition, dependency outcome and attempt count. Add distributed tracing for latency paths and service-level indicators for outcomes customers actually feel.
Observability should answer three questions quickly: what is affected, where did it stop progressing, and what action is safe now?
8. Recovery that depends on memory
A production system eventually encounters work that cannot complete automatically: a provider returns an ambiguous outcome, a message exhausts retries, a database write fails after an external side effect or an operator makes a mistake.
Solution: reconciliation and runbooks
Build reconciliation jobs that compare your internal state with trusted external records. Let operators safely replay a message, correct a state with an audit trail or request a provider status check. Define runbooks before the incident: who owns the decision, which actions are reversible and what evidence is required before changing customer-visible state.
Reliable systems are not systems that never need intervention. They are systems that make intervention bounded, visible and safe.
A practical AWS implementation
On AWS, a typical shape might use an Application Load Balancer and stateless application service for requests, an idempotency table in DynamoDB or a relational database, Amazon SQS for durable asynchronous work, worker services for fulfilment, EventBridge for selected domain events, CloudWatch metrics and alarms, and a dead-letter queue for repeated failures.
The services are not the architecture. The architecture is the contract between them: where intent is stored, how an event can be replayed, who owns the state transition and how recovery is verified. AWS reliability guidance likewise frames resilience as strong foundations, resilient architecture, controlled change and proven recovery processes. Explore the AWS Reliability Pillar.
When not to distribute
Do not add a service boundary merely because a diagram looks modern. A database transaction inside one well-structured application is easier to reason about than a saga across five services. Start with the smallest architecture that provides the required availability, ownership and scaling characteristics. Add distribution only when the failure domain or capacity boundary is real.
The engineering principle to keep
Every distributed operation needs an answer to four questions: Can this work be repeated? Can it arrive late? Who is allowed to decide the state? How do we know it finished?
If those answers are explicit, failures become events the system can absorb. If they are implicit, an ordinary timeout can become a business incident.
Comlabs designs and operates cloud systems that have to keep working after the first clean demo. Explore our AWS Cloud and DevOps services, read our AWS distributed systems guide, or see how we approached application performance and scaling.
