Blog

Failure Injection for Kafka: A Weekend Exercise That Prevents Quarter-End Outages

Table of Contents

Table of Contents

Picture the postmortem. The quarter closes on Friday, and three incidents arrive in one day. A checkout pipeline starts timing out because a broker restart took longer than the team estimated. A disk fills on one broker and turns a routine retention problem into a produce outage. An availability-zone network blip happens to arrive while a leader election was already in progress, and customers see a burst of duplicate and failed requests. By Monday the team has a clear timeline and three root causes, none of which is new.

That is the uncomfortable part. Nothing in that postmortem was unknown. The team had replicas, retries, monitoring, and a runbook. What it did not have was practice under load, so three known failure modes composed into a quarter-end outage. Failure injection for Kafka exists to remove that gap: it proves, before the business peak, that each failure fails the way your configuration says it should, then raises the pass line until the surface you actually run is covered.

You do not need a chaos engineering platform to get there. Four well-designed experiments, run against brokers, disks, zones, and clients in one weekend, will teach you more about your stack than the last quarter of dashboards. And once the habit exists, the same four fields transfer cleanly when a storage architecture change moves the injection points.

1Break it on purpose before the quarter does

Teams postpone failure injection for reasons that sound reasonable in a planning meeting:

  • Production is busy, and the test window is small.
  • The monitoring is probably good enough already.
  • Someone documented the failover behavior, so the stress test feels redundant.

Each of those reasons is a bet that the unknown failure will not arrive at the worst possible time, and the quarter does not respect that bet. Quarter-end load is when retry storms, replay backlogs, and rush-hour deploys all happen at once. If your broker failover has a behavior you did not expect, that is the day it shows up. Injection moves that discovery from the worst day to a Saturday morning when the blast radius is under your control and the people who own the runbook are in the same room.

Treat failure injection as configuration verification, not performance theater. When you set replica placement, min.insync.replicas, producer acknowledgement, and leader election policy, you are making claims about how the system degrades. An injection experiment is how you check those claims outside an incident. It does not tell you the system is perfect; it tells you which of your beliefs were actually tested.

A reasonable sequence is to rehearse each fault in a stage environment that mirrors the production configuration, then confirm the two or three highest-risk cases in production during a declared low-risk window. Some behaviors only appear with real instance types, real network segmentation, and actual on-call engineers. The goal is not to break production on purpose; it is to stop treating the largest blast radius as the one place you are not allowed to look.

2The failure set worth your weekend

Four failure modes cover most of the quarter-end incidents that show up in Kafka postmortems. Start with these, then add topics that match your own incident history.

  • Broker process failure. Stop or kill one broker, then watch leader election, controller behavior, in-sync replica movement, and the producer error rate. This is the least expensive experiment and usually the one with the most stale assumptions about recovery time.
  • Disk or local storage fault. Fill or slow a broker's log disk. Watch whether the cluster isolates the failing broker quickly, how replication behaves, and whether a full disk becomes a produce outage or stays a localized degradation.
  • Availability-zone impairment. Block traffic into one zone. Watch client failover, controller stability, and how much cross-zone traffic the recovery path creates while the rest of the system is serving the business.
  • Client flood and retry storm. Drive producers and consumers above their normal rate, or hold consumers offline so lag accumulates, then let them resume. Watch backpressure, throttling, and whether retry behavior protects the cluster or feeds the storm.

These four were chosen because each exercises a different contract. Broker failure tests the election and metadata path. Disk failure tests the durability and replication path. Zone failure tests the network and placement assumptions. Client flood tests the admission-control path that most teams never configured consciously. Apache Kafka's replication documentation covers the mechanisms; your job is to verify that your specific configuration implements them the way you expect.

Three of the four are also where a passing test can hide a fragile system. A broker restart in an idle cluster proves little about a broker restart during a write-heavy window with active consumer lag. A zone isolation that never triggers leader movement because partitions are pinned proves less than one that does. The fault you inject matters less than the pressure that is present while you inject it.

3Minimal experiment design: hypothesis, injection, pass line

A useful injection experiment has four fields, and most teams only fill in two of them.

  • Hypothesis. State what you believe will happen and why. A weak hypothesis is "the cluster stays up." A useful one is: producers using acks=all will keep seeing bounded retriable errors while one broker restarts, and no committed record will be unavailable, because min.insync.replicas keeps the write path on replicas that remain in sync.
  • Injection. Name the action, the scope, and who can stop it. "Kill broker 3" is not enough. "Stop the Kafka process on broker 3, in zone B, during the 14:00 window, with the on-call engineer holding the rollback command" is an injection.
  • Observation. List the metrics that prove or disprove the hypothesis: produce and fetch latency, UnderReplicatedPartitions, OfflinePartitionsCount, leader-change events, consumer lag, and error-code distribution. Favor signals that connect the failure to customer-visible behavior, not broker liveness alone.
  • Pass line. State the measurable bar in advance. A pass line is a contract: no produce request fails permanently, lag returns to baseline within a duration you set from history, and no topic drops below its target in-sync replica count. Build the number from your own baseline rather than copying a threshold, because a pass line that ignores your workload is decoration.

The fields form a loop. A hypothesis you cannot observe is a press release. An injection without a pass line is a gray exercise where everyone agrees the drill "went fine" and nobody knows what would have failed it. Writing the four fields first forces the team to name the belief being tested, which is the part of an injection that transfers to the next quarter.

Failure injection plan board linking each failure mode to its hypothesis, injection, and pass line

The minimum in-sync replicas documentation is the right place to confirm what the broker-side contract means before you encode it in a pass line. The pass line itself still belongs to your workload and your customers.

4A four-hour Saturday drill

The four experiments compress into a morning if you prepare in advance. This is an illustrative schedule: shift the blocks to match your change window, and expect the first run to take longer because runbook gaps show up while the clock is running.

Weekend drill schedule from Friday preparation through Saturday drill blocks to Sunday production confirmation

The Saturday morning itself breaks into five blocks:

Time blockInjectionWatchPass line
09:00-09:40Stop one broker cleanlyLeader election, ISR churn, producer error codesNo permanent produce failure; lag returns to baseline
09:50-10:30Fill or slow one broker's log diskReplica movement, controller decisions, append latencyCluster isolates the broker; produce path stays available
10:40-11:20Flood producers; hold consumers, then resumeBackpressure, throttling, rebalance timeRetry behavior does not amplify the storm
11:30-12:10Isolate one availability zoneClient failover, controller stability, cross-zone trafficFailover converges; recovery traffic does not trip another alert
12:10-12:40Debrief and decision logFindings against each hypothesisEvery row has a yes, a no, or a changed config

The order is deliberate. Broker failure is the most contained and the least expensive to investigate mid-drill. Disk failure layers the storage path on top of the process path from the previous block. The client flood introduces outside pressure, and the zone isolation combines everything at the largest blast radius. Rearranging the order is fine; skipping the progression from contained to broad makes the broad experiments harder to read.

Before every block, check the rollback. The abort condition should be as specific as the injection: a metric threshold, a business error rate, or a stated recovery failure. If nobody can say the words that stop the test, the test is a live-fire exercise with extra steps.

A failover you have not rehearsed is a failover you will learn about from customers, usually in the worst possible week.

Run the same drill against stage first. Production confirmation is for the cases where stage cannot reproduce the real constraint: real instance types, real network segmentation, and an operations team that is actually on call instead of watching a demo.

5New injection points on shared storage

The four drills above assume a cluster where brokers hold durable partition data on local disks. In that model, a broker failure is partly a process problem and partly a data-placement problem: the cluster has to elect new leaders and copy replicas to recover the replication factor. That coupling is why broker replacement and disk recovery dominate so many postmortems.

A Shared Storage architecture changes which part of the system is the durability boundary. Durable stream data lives in object storage, behind a storage engine such as S3Stream, and the broker becomes a stateless compute node with a fast write-ahead path (WAL storage) and a read cache (Data caching). Apache Kafka's own remote log storage is a related idea: tiered storage moves cold data to remote storage while hot data and replication stay local, so the ISR and local-disk behaviors remain. Shared Storage is a different operating model because the broker stops owning durable data, and the injection points move accordingly.

What a good failure model needs in this world is a clear durability boundary in object storage, a small write path that can survive a broker replacement, and brokers that re-attach without a data rebalance. AutoMQ, a Kafka-compatible streaming platform, builds that shape with S3Stream for durable stream data, WAL storage for the low-latency write path, and Data caching for hot reads. Once brokers are stateless, the experiment design changes in useful ways.

  • Broker restart becomes a re-attach test. With no durable partition data on the broker, the recovery step stops being "copy replicas to a new node." The experiment measures how fast the replacement rejoins the shared stream and resumes service, and whether consumers see a bounded gap.
  • Object-storage availability becomes a first-class fault. Throttle or deny the object-storage path and watch the WAL upload backlog and producer acknowledgement behavior. The durability boundary has moved, so this is where a loss would appear first.
  • WAL storage becomes its own failure domain. Isolate the WAL volume and watch the failover path and whether low-latency writes stay bounded. A broker can be disposable while its WAL is still a critical local dependency.
  • Broker-to-object-storage partitions replace some broker-to-broker partitions. Partition brokers away from the object-storage endpoint and watch produce latency, upload lag, and read-path behavior, separate from the inter-broker metadata path.

Injection point map for a shared storage architecture, highlighting object storage, WAL, and broker re-attach failures

On this list, the cross-AZ accounting also changes. With shared object storage, replication does not shuttle copies between brokers across zones, so a zone isolation test should confirm that the recovery path does not quietly reintroduce cross-AZ data transfer. That assumption is worth a dedicated experiment, not a footnote in the zone drill.

These are the same four fields with different content: the hypothesis, injection, observation, and pass line now describe object storage, WAL, and broker re-attach instead of local-disk replicas. Run that set against any shared-storage deployment you are evaluating, AutoMQ included. The compatibility documentation is the place to confirm which client and protocol behaviors carry over unchanged if you are comparing against your current stack.

Go back to the quarter-end postmortem. Three incidents, all known, all unpracticed. The question is not whether your stack could survive them; it is whether you have evidence, captured on a Saturday, that it will. If the recurring answer points at broker-local storage and the data movement around it, run the same drill against a shared-storage deployment and watch the failure model move with it. The AutoMQ open-source project is a practical place to start the shared-storage rows of that table.

6References

7FAQ

7.1Do I need a chaos engineering platform for Kafka failure injection?

No. The value comes from stating a hypothesis, injecting one fault, and measuring against a pass line you set from your own baseline. Four weekend experiments against brokers, disks, zones, and clients cover more ground than a platform dashboard without a pass line. A dedicated tool becomes worth it when you want the injections automated and repeated on a schedule, not when you start.

7.2Should failure injection run in production or only in staging?

Start in a stage environment that mirrors production configuration, then confirm the highest-risk cases in production during a declared low-risk window. Some behaviors, such as real network segmentation, instance-type behavior, and an actual on-call response, do not reproduce in staging. Production injections need a specific rollback command, an abort condition, and an owner named before the test begins.

7.3What is the pass line for a Kafka failover test?

The pass line is a workload-specific contract: no produce request fails permanently, committed records remain available, lag returns to baseline within a duration you set from history, and no topic drops below its target in-sync replica count. Write the line before the injection, because a test that runs first and sets the bar afterward cannot fail honestly.

7.4Does shared storage remove the need for Kafka failure injection?

No. It moves the injection points. Broker restart becomes a re-attach test, object-storage availability and WAL storage become first-class faults, and broker-to-object-storage partitions replace some broker-to-broker partitions. The four fields stay the same; the failure set and the durability boundary change.

Newsletter

Subscribe for the latest on cloud-native streaming data infrastructure, product launches, technical insights, and efficiency optimizations from the AutoMQ team.

Join developers worldwide who leverage AutoMQ's Apache 2.0 licensed platform to simplify streaming data infra. No spam, just actionable content.

I'm not a robot
reCAPTCHA

Never submit confidential or sensitive data (API keys, passwords, credit card numbers, or personal identification information) through this form.