10 min readArchitecture
Kafka Is Back. That's the Dangerous Part.
The brokers are healthy again. The dashboards are green. Someone in the incident channel types “we’re back” and the tension drops.
Then the publisher restarts, a million pending events go out as fast as the network allows, and core banking — which comfortably handles two thousand transactions per second — receives several times that within a minute. Response times go from 150ms to seconds. Consumers hit their timeouts and retry, which adds load, which pushes latency higher, which causes more timeouts.
The second incident is now worse than the first, and unlike the first one, you caused it.
If you have built the transactional outbox, you already did the hard part: nothing was lost while Kafka was down. Events accumulated safely in a database table exactly as designed. But durability guarantees only cover the storing. What happens when you turn the tap back on is a separate design problem, and it is usually the one nobody has thought about.
The arithmetic nobody does during the incident
A million events. Downstream sustains 2,000 TPS. At full capacity that backlog takes over eight minutes to clear — and that is the optimistic figure, because it assumes downstream has nothing else to do.
It does. Live traffic did not stop. Customers are still paying. If the platform normally runs at 1,200 TPS, the actual headroom for backlog is 800 TPS, and the drain takes over twenty minutes. Push harder than that and you are not draining faster, you are queueing inside somebody else’s system.
That is the part that makes recovery dangerous: Kafka’s throughput is not the constraint. Kafka will happily accept a million events in seconds. So will your relay. The constraint lives two hops downstream, in a core banking system that was sized for steady state, and nothing between here and there knows that.
Retries make it worse rather than better. A timeout at 2,000 TPS produces a retry, so the offered load becomes 2,000 plus retries. Latency rises, more calls time out, more retries. This is a feedback loop with positive gain — the textbook definition of an unstable system — and it will not settle on its own.
A fixed drain rate is a guess
The common advice is to rate-limit the replay: pick a conservative number, 500 or 800 TPS, and publish at that.
That is much better than flushing. It is also a number someone invented, and it is wrong in both directions.
It is too high when the outage happened at peak hour and downstream is already near capacity. It is too low at 3am when there is spare capacity and you are needlessly extending the recovery — every minute of backlog is a customer wondering where their money went.
And downstream capacity is not a constant anyway. It varies with time of day, with what batch jobs are running, with whatever else the bank has going on. A constant cannot track a variable.
The fix is to stop guessing and close the loop: let the drain rate be an output of downstream health rather than a configuration value. This is the same idea as TCP congestion control — increase gently while things are healthy, back off sharply the moment they are not.
/**
* AIMD drain governor: additive increase, multiplicative decrease.
*
* Ramps up slowly while downstream stays healthy and backs off hard the moment
* it does not. The asymmetry is deliberate — recovering capacity should be
* cautious, shedding load should be immediate.
*/
@Component
public class DrainGovernor {
private static final int MIN_TPS = 50;
private static final int MAX_TPS = 3_000;
private static final int INCREASE_STEP = 50; // additive
private static final double DECREASE_FACTOR = 0.5; // multiplicative
private final AtomicInteger targetTps = new AtomicInteger(MIN_TPS);
private final DownstreamHealth health;
/** Re-evaluate once per second; slower than that and the loop lags reality. */
@Scheduled(fixedRate = 1_000)
public void adjust() {
if (health.isDegraded()) {
int reduced = Math.max(MIN_TPS,
(int) (targetTps.get() * DECREASE_FACTOR));
int previous = targetTps.getAndSet(reduced);
if (reduced < previous) {
log.warn("Downstream degraded — drain rate {} → {} TPS", previous, reduced);
}
} else {
targetTps.updateAndGet(current -> Math.min(MAX_TPS, current + INCREASE_STEP));
}
}
public int currentTps() {
return targetTps.get();
}
}
The health signal is where the judgement lives. Success/failure is too blunt — by the time calls are failing, you are already in trouble. Latency degrades before errors appear, which makes it the earlier and more useful signal:
@Component
public class DownstreamHealth {
private static final Duration LATENCY_CEILING = Duration.ofMillis(400);
private static final double ERROR_CEILING = 0.02;
private final MeterRegistry metrics;
boolean isDegraded() {
double p99 = metrics.timer("bank.call").takeSnapshot()
.percentileValue(0.99, TimeUnit.MILLISECONDS);
double errorRate = errorRateOverLastSeconds(10);
// Latency moves first; errors are the lagging confirmation.
return p99 > LATENCY_CEILING.toMillis() || errorRate > ERROR_CEILING;
}
}
Set LATENCY_CEILING from the healthy baseline, not from the SLA. If core banking normally answers in 150ms, then 400ms already means something is wrong — waiting for it to breach a 2-second SLA means you back off long after the queues have built.
The relay then asks the governor how fast it may go, instead of running flat out:
@Scheduled(fixedDelay = 100)
public void drainBacklog() {
int budget = governor.currentTps() / 10; // this tick's share of the second
if (budget <= 0) return;
List<OutboxRow> rows = outbox.claimBatch(budget);
if (rows.isEmpty()) return;
publish(rows);
backlogMetrics.recordDrained(rows.size());
}
This changes the operational posture completely. Nobody has to sit there watching a dashboard and adjusting a config value. The system finds the fastest rate downstream can actually absorb, continuously, and gives it up the instant that changes.
Live traffic must beat the backlog
Here is a decision that gets made by accident far more often than deliberately.
When the backlog is draining, two kinds of work compete for the same downstream capacity: new payments happening right now, and events from two hours ago. If they share a queue, they share a fate — and new payments start queueing behind old ones.
That is the wrong priority. A customer waiting on a payment right now is watching a spinner. An event from two hours ago belongs to a customer who has already seen a timeout, already worried, possibly already retried. Making live traffic wait behind the backlog turns a resolved incident into a fresh one for everybody.
So separate the lanes, and let live traffic starve the backlog when it needs to:
public int backlogBudget() {
int downstreamCapacity = governor.currentTps();
int liveDemand = liveTrafficMeter.currentTps();
// Live traffic is served first; the backlog gets what is left, and it is
// allowed to get nothing at all.
return Math.max(0, downstreamCapacity - liveDemand);
}
Math.max(0, ...) is doing something important: at peak load the backlog drain rate can legitimately fall to zero. That feels wrong the first time you see it — the backlog stops shrinking during the busiest hour — but it is correct. The backlog is durable and patient. It will still be there at 9pm when there is capacity, and nothing is lost by waiting.
Not everything in the backlog is still worth sending
This is the part most recovery designs miss entirely, and it is not a technical problem.
An event that has been sitting for two hours may no longer describe reality. The customer saw a timeout and retried, so there is now a second, successful transaction. Or the session expired. Or a reconciliation job already settled it through another path. Replaying it blindly is technically correct and commercially wrong — and idempotency will not save you here, because these are genuinely different transaction IDs representing the same customer intent.
So the backlog needs triage before it needs throughput:
public DrainDecision triage(OutboxRow row) {
Duration age = Duration.between(row.createdAt(), Instant.now());
if (age.compareTo(Duration.ofHours(24)) > 0) {
// Too old to replay safely into a live system.
return DrainDecision.reconcile("Exceeded replay window");
}
if (paymentLedger.hasSettledEquivalent(row.transactionId())) {
// Already resolved through retry or another route.
return DrainDecision.discard("Already settled");
}
return DrainDecision.publish();
}
reconcile is not discard. Those events go to a separate topic for a reconciliation process — and often a human — to resolve. They are not thrown away, they are just not fired into a live payment path hours after the fact.
The replay window is a business decision, not an engineering one. Ask your payments and operations people what the cutoff is; they will have a clear view, and it may well be far shorter than 24 hours.
The signals that actually matter
Backlog size is the metric everyone puts on the dashboard, and it is the least useful of the set. A backlog of 400,000 that is shrinking steadily is fine. A backlog of 5,000 that has not moved in ten minutes is an outage.
The ones worth alerting on:
- Age of the oldest unpublished event. This is the real customer-impact metric. Size tells you about volume; age tells you how long someone has been waiting.
- Drain rate against arrival rate. If arrivals exceed drain, the backlog is growing and you will never catch up — that is a different problem needing more capacity, not more patience.
- Estimated time to clear, derived from the two above. This is the number the incident channel actually wants, and it is the one people otherwise compute badly in their heads.
- Downstream p99 during the drain. The whole point of the control loop, and the thing that tells you it is working.
And the combination worth watching hardest: backlog shrinking while downstream latency climbs. That means the drain is winning at downstream’s expense. A closed-loop governor handles it automatically. If you are draining at a fixed rate, that pattern is your cue to intervene, and it will show up well before anything starts failing.
Ordering and duplicates still apply
Two things carry over from the steady-state design and do not change during recovery, so I will keep this short.
Key by account, not transaction. Two operations on the same account must not be reordered during a burst replay — a debit and a credit applied out of order can produce a transient negative balance that trips fraud rules or bounces a subsequent payment. Keying by account puts them on the same partition, in order.
Idempotency is mandatory, not optional. Replay means at-least-once, and duplicates are certain rather than theoretical. Enforce it with a unique constraint in the same transaction as the work, never with a SELECT-then-if. I covered the mechanics in the outbox post.
Rehearse it, or it does not work
Recovery code has an unusual property: it only ever runs during an incident. Which means, unless you do something about it, it is exercised for the first time at the worst possible moment, by people under pressure, in the middle of the night.
It will not work. Someone will have refactored the relay and broken the governor. The MIN_TPS will have been set to something absurd during a debugging session and never reverted. The triage rule will reference a table that was renamed.
The fix is unglamorous: pause the publisher deliberately in a non-production environment, let a realistic backlog build, and turn it back on while watching the control loop respond. Do it quarterly. You will find something every time.
The teams that recover well from outages are not the ones with the cleverest architecture. They are the ones who have done it before on a Tuesday afternoon with nothing at stake.
The short version
- The outage is the easy part. The outbox already saved your data. Recovery is where the damage happens.
- Do the arithmetic: backlog ÷ (downstream capacity − live traffic). That is your real drain time.
- A fixed drain rate is a guess. Close the loop and let downstream latency set the pace.
- Use latency as the signal, not errors. Errors are the lagging indicator.
- Live traffic wins. The backlog can legitimately drain at zero during peak.
- Triage before publishing. Some events are too old to replay and belong in reconciliation.
- Alert on backlog age, not backlog size.
- Rehearse it, or you are debugging it live.
Kafka being up does not mean the system is healthy. The goal after an outage is not to clear the backlog quickly — it is to clear it without causing the next one.