Handling Partial Failures in Distributed Systems

A payment either happens or it doesn't. That's the story we tell customers, and it's the story regulators want to hear. But underneath, a single card...

Originally published onanselmfowel.com

A payment either happens or it doesn't. That's the story we tell customers, and it's the story regulators want to hear. But underneath, a single card authorization might touch six services, two message queues, a fraud model, and a ledger, and any one of those can time out while the others keep going. The interesting failures in distributed systems are almost never the clean ones where everything is down. They're the ones where half the work succeeded.

I've spent most of my career on the plumbing behind money movement, and if I had to name the thing that separates teams who sleep at night from teams who don't, it's how seriously they take partial failure. Not total outages. Those are loud and honest. Partial failures are quiet, and they lie to you.

The lie of the happy path

Early in my career I inherited a service that captured payments and then, in the same request, called a downstream service to update a merchant's available balance. The code assumed both would work. When the balance service was slow, the capture would time out on the client side, the caller would retry, and we'd occasionally capture the same payment twice. Not often. Maybe one in forty thousand transactions. Which sounds fine until you remember that at a few million transactions a day, that's dozens of double-charges daily, each one a support ticket and a small dent in trust.

The root cause wasn't a bug in the usual sense. Every line of code did exactly what it said. The problem was that the design treated "the request completed" and "the work completed" as the same fact, and in a distributed system they are two different facts. A timeout tells you that you stopped waiting. It tells you nothing about whether the other side finished.

The most dangerous state in any distributed system is "I don't know." Down is safe, because you retry. Up is safe, because you move on. "I sent the request and never heard back" is the state that generates the 2am pages and the awkward reconciliation meetings.

Idempotency is not optional

The single highest-leverage habit I've pushed on every team I've run is that any operation which changes money, or anything that feels like money, must be idempotent. The caller supplies a key, the server remembers what it did with that key, and repeating the request returns the original outcome instead of doing the work twice. This is not exotic. Stripe has done it publicly for years, and it's the backbone of why you can safely retry a charge without sweating.

The part people get wrong is the storage. An idempotency key that lives in memory, or in a cache with a short TTL, is worse than none at all, because it gives you false confidence. The record of "I processed key X and here's the result" has to be as durable as the money movement itself, ideally in the same transaction. Here's roughly the shape I want to see:

BEGIN TRANSACTION;

-- Claim the key. If it already exists, we've done this work.
INSERT INTO idempotency_keys (key, status, created_at)
VALUES (@key, 'in_progress', SYSUTCDATETIME());
-- A unique constraint on key makes the duplicate INSERT fail,
-- which is exactly what we want: the second caller loses the race.

INSERT INTO payments (id, merchant_id, amount_cents, currency, idempotency_key)
VALUES (@paymentId, @merchantId, @amount, @currency, @key);

UPDATE idempotency_keys
SET status = 'succeeded', response_payment_id = @paymentId
WHERE key = @key;

COMMIT;

The second caller hits the unique constraint on the key, the transaction rolls back, and the API layer reads the stored result and returns it. No double payment. The trade-off is real: you're adding a write and a constraint check to your hot path, and you need a cleanup job for stale in-progress rows. I'll take that cost every time over explaining to a merchant why they were debited twice.

Retries that don't make it worse

Retries are the first tool everyone reaches for, and they're also the most common way teams turn a small problem into an outage. A downstream service gets slow, every caller retries three times, and now that service is handling four times its normal load precisely when it's least able to. I've watched a minor database blip cascade into a full regional failure this way, and the retries were the fuel.

If you're going to retry, retry with discipline. That means a few non-negotiable rules:

  • Exponential backoff with jitter, so a thousand clients don't all retry in the same 50ms window and hammer the recovering service in synchronized waves.
  • A hard cap on total attempts and total elapsed time. Retrying for two minutes on a request the user abandoned after five seconds helps no one.
  • Retry only on failures that might actually be transient: timeouts, connection resets, 503s. A 400 or a validation error will fail identically forever, so retrying it is just noise.
  • An idempotency key on every retried write, per the previous section, so the retry can't duplicate work.

My rule of thumb: if you can't explain why a specific retry is safe, it isn't. Retries and idempotency are two halves of the same design decision. Doing one without the other is how you get the double-capture story I opened with.

Circuit breakers and knowing when to stop

There's a point where retrying stops being optimism and becomes denial. Circuit breakers exist for that point. When a dependency has failed enough times in a short window, you stop calling it, fail fast for a cooldown, then let a trickle of requests through to test whether it's healthy again. The breaker protects both sides: your threads stop piling up waiting on a dead service, and the dead service gets breathing room to recover.

In .NET I lean on Polly for this, wired through the standard resilience pipeline so the policy lives in one place rather than scattered across call sites. The configuration matters more than the library. A breaker that trips too eagerly turns a two-second hiccup into a thirty-second outage for your customers. One that trips too late doesn't protect anything. I usually start conservative, watch the real failure distribution for a week, and tune from there. Guessing the thresholds up front is how you end up with a breaker that either never fires or fires constantly.

The dual-write trap

Here's the pattern that has burned me more than any other: write to the database, then publish an event to a queue. Two systems, one logical operation, no shared transaction. If the database commit succeeds and the publish fails, you have a payment that exists but that nothing downstream knows about. If you flip the order, you can publish an event for a payment that never actually got saved. There is no ordering of two independent writes that is safe, and yet I still see this pattern shipped constantly because it reads as perfectly reasonable in a code review.

The fix that actually holds up under load is the transactional outbox. You write the business change and an "outbox" row in the same local transaction, so they're atomic. A separate process reads unpublished outbox rows and pushes them to the queue, marking them sent once the broker acknowledges. Because delivery is now at-least-once, your consumers have to be idempotent, which loops right back to the discipline from earlier. It's more moving parts than a naive publish, and I won't pretend it's free. But it converts "silently lost payment event" into "event arrives, possibly more than once," and the second problem is one you can actually solve.

Timeouts are a product decision

Most timeout values I encounter in the wild were never chosen. They're whatever the HTTP client defaults to, often 100 seconds or infinity, and nobody touched them. That's a mistake, because a timeout is really a statement about how long a piece of work is allowed to matter. Holding a card network connection open for 100 seconds isn't patience, it's a slow-motion thread pool exhaustion waiting to happen.

I want timeouts set deliberately at every hop, and I want them to add up to something shorter than the caller's patience, not longer. If a user's browser gives up after 30 seconds, there is no point in an internal service waiting 60 to hear back from its dependency. Set the inner timeouts so the whole chain fails before the outer one does, and propagate a deadline down the call graph so a request that's already doomed doesn't keep consuming resources. A request nobody is waiting for anymore is pure waste, and under load that waste is exactly what pushes you over the edge.

Reconciliation is the safety net

You will not prevent every partial failure. Accept that early. Networks partition, processes die between the commit and the acknowledgment, a deploy restarts a pod mid-flight. The question isn't whether inconsistency happens, it's how fast you notice and how cheaply you fix it. That's what reconciliation is for, and in regulated fintech it's not optional anyway, because someone will eventually ask you to prove that your internal ledger matches the money that actually moved.

The teams I trust run automated reconciliation continuously, not as a quarterly fire drill. A job walks your records against the payment processor's records and against your own ledger, flags anything that doesn't line up, and either auto-resolves the known-safe cases or parks the rest for a human. The first time you run one of these against a mature system, it will find things. It always does. I've never once seen a first reconciliation come back perfectly clean, and the ones that look clean usually mean the reconciliation itself is broken.

Observability for the in-between

You can't fix what you can't see, and partial failures are specifically the things that don't show up in a simple up/down dashboard. Both services are green. The queue depth looks normal. And yet a slice of transactions is stuck in limbo. To catch that, you have to instrument the states between "started" and "done," not just the endpoints.

Concretely, I want a metric on how many operations are currently in-flight past their expected duration, an alert on outbox rows that haven't been published within N seconds, and correlation IDs that let me trace a single payment across every service it touched. The correlation ID is the cheapest of these and the one I'd fight hardest to keep. The first time you're staring at a stuck payment at 2am and can paste one ID into your logs to see every hop it made, you'll understand why. Without it, you're grepping timestamps across a dozen services and guessing.

Anselm Fowel, CTO and fintech architect
Anselm Fowel — CTO & fintech architect

Conclusion

If there's one shift I'd push on any engineer moving into this world, it's this: stop designing for the case where things work, and stop designing only for the case where things are fully down. Design for the ugly middle, where the request half-finished and you're not sure. That's where the money gets lost and the trust gets spent. Assume every call can fail after doing its work but before telling you, and build so that a retry, a duplicate, or a crash at the worst possible moment is boring instead of catastrophic. Boring is the goal. In systems that move money, the highest compliment you can pay an architecture is that its failures are dull, well-understood, and already handled before anyone had to wake up for them.

Chat with us