Partial Failure, Ordering, and Eventual Consistency
Understand partial API failures, message ordering, and eventual consistency, then design status models and recovery paths clients can reason about.
Partial failures leave clients unsure whether a timed-out request committed, while downstream views may lag behind the authoritative service. This guide shows how to expose durable operation states, apply per-entity event versions, and choose bounded retries, saga compensation, or operator repair. It also covers reconciliation, observability, and recovery-capacity estimates so teams can detect stale projections and restore workflows safely.
Partial Failure, Ordering, and Eventual Consistency
Introduction
A checkout request can reach the server, reserve inventory, and then time out before the response reaches the browser. The client cannot tell whether the server did nothing or committed work, so blindly retrying may create a second workflow.
POST /checkouts
Idempotency-Key: checkout-123
A stable key helps prevent duplicate effects, but the client still needs a way to learn whether the operation is pending, complete, or being compensated. This guide covers status models, per-entity ordering, and recovery paths for delayed or missing updates.
Model the outcome explicitly
Return acceptance separately from completion when work continues asynchronously. A response can give the caller a stable operation ID and status URL:
{
"operationId": "op_7f32",
"status": "processing",
"statusUrl": "/operations/op_7f32"
}
Define what each state means. accepted means the service recorded the request; processing means work is underway; succeeded and failed are terminal outcomes; compensating means the system is undoing or offsetting earlier steps. A client can poll the status resource with bounded intervals or receive a notification. Include a version or update time so clients can recognize stale status, and avoid reporting succeeded until the operation’s promised effects are complete.
When to use eventual consistency
It fits workflows where independent services own their data, where read models can lag briefly, or where asynchronous processing protects availability. It is a poor fit for an invariant that must be checked atomically, such as preventing a bank account from spending the same balance twice. Keep strong consistency for the narrow operation that protects the invariant, then distribute resulting events afterward.
| Approach | Latency | Availability | Stale-read behavior | Complexity | Recovery |
|---|---|---|---|---|---|
| Strongly consistent operation | Higher when coordination is required | A participant or quorum failure can block writes | Reads reflect committed writes within the consistency boundary | Coordination and contention | Retry the operation safely; resolve failures inside its transaction boundary |
| Eventual consistency | Lower for accepted asynchronous work | Services can accept work while another service is unavailable | Read models can lag; expose versions or freshness | Pending states, ordering, and reconciliation | Replay or reconcile missed updates |
| Bounded retry | Adds delay while waiting for a dependency | Preserves the workflow if the dependency recovers | Existing projections may remain stale during retries | Retry limits, backoff, and idempotency | Retry while the step is still valid; stop at a defined limit |
| Saga compensation | Adds time for compensating steps to complete | Other participants can continue independently | Intermediate effects remain visible until compensation finishes | Each compensation needs its own state and failure handling | Apply a new business action, then reconcile or route failures for repair |
Compensation is not rollback
In a workflow spanning several services, a later failure cannot usually undo earlier commits with one database rollback. A saga handles this by issuing compensating actions, such as releasing a reservation after payment fails. Compensation is a new business operation: it can fail, take time, or produce a result that differs from the original state. For example, a refund may reverse a charge, but it does not erase the charge from the ledger.
Choose between retrying the failed step and compensating based on the business rule. If the dependency may recover and the action remains valid, retry with a bounded policy. If the workflow can no longer continue, run an explicit compensation. Make each step and compensation safe to repeat, record its outcome, and expose intermediate states such as compensating or compensation_failed through the operation status API. Provide a manual repair path for cases where compensation itself cannot complete.
flowchart TD
A[Partial failure leaves some steps committed] --> B[Persist operation as pending]
B --> C{Can the failed step still complete?}
C -->|Yes| D[Retry the idempotent step within its limit]
D --> E{Did the retry succeed?}
E -->|Yes| F[Resume the workflow]
E -->|No| G[Compensate completed steps]
C -->|No| G
G --> H{Did compensation succeed?}
H -->|No| I[Queue operator repair]
F --> J[Reconcile authoritative state and projections]
H -->|Yes| J
I --> J
J --> K[Repair gaps and verify entity versions]
K --> L[Publish terminal status and projection freshness]
The reconciliation step checks the final state even when retries or compensation report success. That catches lost notifications and projections that missed an update.
Key Takeaways
- Model compensation as a durable, repeatable operation, not an automatic rollback.
- Keep the original and compensating actions visible in the operation history.
- Decide whether to retry, compensate, or require operator repair for each failure point.
Implementation snippet: version check
A consumer can ignore stale events and flag a sequence gap for repair. The update and stored version should be atomic.
async function apply(event: AccountEvent): Promise<void> {
await db.transaction(async (tx) => {
const current = await tx.getVersion(event.accountId);
if (event.version <= current) return; // duplicate or stale
if (event.version !== current + 1)
throw new Error("sequence gap; retry or reconcile");
await tx.apply(event);
await tx.setVersion(event.accountId, event.version);
});
}
Production failures and mitigations
The event is committed but the API response is lost, so the client creates a duplicate workflow. Use an idempotency key. A consumer crashes after updating its database but before acknowledging the message; make the handler idempotent. The projection lags for minutes while clients poll aggressively; expose processing status and use bounded polling or notifications. A poison event blocks a strict-order partition; quarantine it, alert, and define whether later events may proceed. Reconciliation jobs should compare authoritative state with projections and repair gaps.
Observability checklist
- Measure end-to-end workflow duration and time spent in each state.
- Track consumer lag, retry counts, dead letters, and sequence gaps.
- Compare projection freshness with an explicit lag objective.
- Log operation IDs, entity versions, and event IDs across services.
- Alert on stuck nonterminal operations and reconciliation backlog.
Estimate recovery capacity
Size recovery work from measured arrival and processing rates rather than a generic target. For one example, assume an integration normally receives 100 events per second, processes only 80 per second while a dependency is degraded, and stays degraded for 30 minutes. The queue grows by (100 - 80) × 1,800 = 36,000 events.
If the recovered service can process 180 events per second while new work continues to arrive at 100 events per second, its net catch-up rate is 80 events per second. Draining 36,000 events then takes 36,000 ÷ 80 = 450 seconds, or 7.5 minutes. During that period, a projection behind this queue can remain stale; track its oldest event age and version lag so clients can distinguish stale reads from missing data.
Use effective throughput in the estimate. If retries, database limits, or rate caps reduce the recovery rate to 144 events per second, net catch-up is 44 per second and the same backlog takes about 13.6 minutes to clear. Include reconciliation records and retry attempts in the measured work rate, and reserve capacity for new traffic. Otherwise, catch-up can keep the queue from shrinking even after the dependency recovers.
Security and Compliance Notes
Authorize status lookups against the same tenant boundary as the original request. Do not expose internal event payloads through a public status endpoint. Validate event producer identity and schema. Keep enough operation and event history to support audit and reconciliation, while applying retention limits and redacting sensitive fields from logs and status responses.
Common Pitfalls / Anti-Patterns
A frequent mistake is labeling a command accepted as completed, or returning stale data without identifying it as a projection. Avoid “exactly once” claims: use idempotent effects, sequence checks, and reconciliation instead. Do not retry a non-idempotent side effect blindly after a timeout; the first attempt may already have succeeded.
Quick Recap Checklist
- State which records are authoritative and which views may lag.
- Expose operation status when completion can take time or remain uncertain.
- Define ordering at the entity or partition level and use versions to detect stale events.
- Make consumers replay-safe and provide reconciliation for inconsistent state.
Interview Questions
Further Reading
- Event sourcing explains how to represent domain changes as an append-only event history.
- API idempotency and safe replays covers duplicate requests and replay-safe side effects.
- Retries, timeouts, backoff, and circuit breakers covers bounded retry behavior when dependencies fail.
- Microsoft Azure Architecture Center: Saga pattern — workflow coordination and compensating transactions across services.
- Microsoft Azure Architecture Center: Event Sourcing — event history, projections, and eventual consistency.
Conclusion
Partial failure is normal once an operation crosses service boundaries. Make uncertainty visible through operation states, define ordering at the entity level, and give operators reconciliation tools. Clients can then handle delay without turning every temporary mismatch into a duplicate command.
Category
Related Posts
Ordering Guarantees in Distributed Messaging
Understand how message brokers provide ordering guarantees, from FIFO queues to causal ordering across partitions, and the trade-offs in distributed systems.
Event Envelopes and Metadata for Reliable EDA
Learn how event envelopes separate transport context from payload, apply CloudEvents attributes, propagate correlation IDs, and validate messages safely.
API Clients, Servers, and Network Boundaries Explained
Understand what API clients and servers each own, how network boundaries fail, and how timeouts, retries, and trust boundaries shape reliable integrations.