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.

published: reading time: 9 min read author: GeekWorkBench
Quick Summary

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

1. What is a partial failure?
It is an outcome where some components completed work while another component or the caller did not observe completion. Distributed systems cannot assume a timeout means that nothing happened.
2. How can a consumer preserve event order?
Partition or serialize events by the entity whose order matters, attach a sequence number, and reject or buffer events that arrive with a gap. Global ordering is usually expensive and rarely needed.
3. When is eventual consistency inappropriate?
When a decision must enforce an immediate invariant, such as preventing duplicate spending or overselling a scarce item. Keep that invariant within a strongly consistent boundary, then publish its result asynchronously.
4. How does a saga differ from a database transaction?
A database transaction can roll back uncommitted changes within its boundary. A saga coordinates committed steps across services and uses retries or compensating actions when a later step fails.
5. Why is compensation not the same as rollback?
Compensation is a new business action that may fail or have its own visible effects. A refund can offset a charge, for example, but the ledger still records both entries.
6. When should a workflow retry instead of compensate?
Retry when the dependency may recover and the operation is still valid, using a bounded policy. Compensate when the workflow cannot continue and the business rule calls for reversing earlier effects.
7. What should an API expose while compensation is running?
Return an operation status that distinguishes the original failure from compensation in progress or compensation failure. Include a stable operation ID so clients and operators can check the same workflow.

Further Reading

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.

#distributed-systems #messaging #kafka

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.

#event-driven-architecture #cloudevents #messaging

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.

#api-design #networking #distributed-systems