Testing Kafka Consumers
bankingSeptember 30, 2026

Testing Kafka Consumers

Testcontainers, Embedded Brokers, and What Actually Catches Bugs

A green @EmbeddedKafka test suite is one of the more convincing false negatives I test against. It runs fast, it passes reliably, and the bug it didn't catch — a consumer that silently double-processes a batch of transactions the moment a second pod joins the consumer group — showed up in production during a routine scaling event, not a chaos experiment. This is about why that gap exists and which tests actually close it. 


What an embedded broker actually is, and why "passing" doesn't mean what it looks like 


@EmbeddedKafka from spring-kafka-test runs a real Kafka broker process, in-memory, inside your test JVM. That's genuinely useful for a specific, narrow purpose: proving your @KafkaListener is wired correctly, your serializer doesn't throw on your actual message shape, and your Spring configuration boots without error. It's fast, and it's honest about what it's testing. 


What it isn't is a reliable proxy for how your consumer behaves during a rebalance — and this isn't a vague "integration tests are always weaker than production" caveat, it's a specific, documented limitation. Spring Kafka's own test support disables Kafka's KRaft mode by default in @EmbeddedKafka, explicitly because of instabilities observed running the newer consumer group protocol (KIP-848) under it — and that new protocol, the one an increasing share of production clusters are moving toward, is only fully testable against a real cluster, not the lightweight broker @EmbeddedKafka builds internally. If your production brokers run KRaft with the new group coordinator and your test suite runs an embedded broker that's quietly falling back to the classic protocol, your rebalance tests are passing against a coordination mechanism your production traffic doesn't actually use. 


The bug class this actually misses 


Rebalance-triggered duplicate processing is the specific failure worth designing a test around: a consumer is mid-batch when a rebalance strips its partition assignment away — a second instance scaled up, a pod restarted, a deploy rolled — and if the offset for that batch wasn't committed before the partition was revoked, the newly-assigned consumer picks up from the last committed offset and reprocesses records the first consumer already handled. On a payment or transaction pipeline, "reprocessed" isn't a log noise problem — it's the exact duplicate-processing scenario idempotency keys exist to catch downstream, and a test suite that never actually triggers a real rebalance will never exercise that path at all. 


 1 @Testcontainers 
 2 class ConsumerRebalanceIntegrationTest { 
 3  
 4     @Container 
 5     static KafkaContainer kafka = new KafkaContainer(DockerImageName.parse("apache/kafka:3.8.0")); 
 6  
 7     @Test 
 8     void reassignedPartitionResumesFromLastCommittedOffset() throws Exception { 
 9         String topic = "tx-events-" + UUID.randomUUID(); // unique per test — no shared state 
10         String groupId = "settlement-test-" + UUID.randomUUID(); 
11  
12         try (KafkaProducer producer = producer(kafka.getBootstrapServers())) { 
13             produceRecords(producer, topic, 100); 
14  
15             try (KafkaConsumer consumerA = consumer(kafka.getBootstrapServers(), groupId)) { 
16                 consumerA.subscribe(List.of(topic)); 
17                 ConsumerRecords firstBatch = consumerA.poll(Duration.ofSeconds(5)); 
18                 // deliberately do NOT commit — simulates a crash mid-batch 
19  
20                 try (KafkaConsumer consumerB = consumer(kafka.getBootstrapServers(), groupId)) { 
21                     consumerB.subscribe(List.of(topic)); 
22                     // joining triggers a real rebalance against a real group coordinator 
23                     ConsumerRecords resumed = pollUntil(consumerB, 
24                         records -> !records.isEmpty(), Duration.ofSeconds(15)); 
25  
26                     assertThat(resumed.count()).isEqualTo(firstBatch.count()); // full re-delivery, not partial 
27                 } 
28             } 
29         } 
30     } 
31 } 

 The assertion worth writing here isn't "a rebalance happened" — that's an implementation detail, not an outcome. It's an invariant: every record from the uncommitted batch gets redelivered exactly once to whichever consumer ends up owning that partition, with none silently dropped and none double-delivered to two active consumers at once. That invariant is what your idempotency-key handling downstream is actually built to tolerate — this test is what proves the tolerance is being exercised, not just theorized about. 


What Testcontainers gets you that a mock or an embedded broker can't 

The reason this test is trustworthy where a mocked Consumer or an embedded broker's simulated coordination wouldn't be: KafkaContainer runs an actual Kafka broker process in Docker, with a real group coordinator, real heartbeats, and real partition-assignment protocol exchange between consumer instances. When consumerB.subscribe() triggers a rebalance in the test above, that's the same coordination machinery — JoinGroup, SyncGroup, partition revocation and reassignment — that runs in production, not a simulated approximation of it. That's what makes the test capable of catching a bug in your offset-commit timing relative to partition revocation, which is precisely the class of bug an in-process embedded broker's simplified coordination model is the least likely to reproduce faithfully. 


What it still doesn't get you 


Worth being honest about, so the test suite doesn't overclaim what it covers: a single KafkaContainer is a real broker, but it's still one broker, not a cluster. It won't reproduce leader movement across brokers, rack-awareness behavior, a rolling broker upgrade, or a group coordinator failing over after a node loss — those need a genuinely multi-broker Testcontainers setup, which is heavier, slower, and worth reserving for the specific tests that need exactly that failure mode, not the default for every consumer test in the suite. 


Test hygiene matters more here than in most integration testing, because Kafka's coordination is asynchronous and stateful by nature: a unique topic and consumer group ID per test avoids one test's leftover group state leaking into another's assertions, and waiting on an observable condition with an explicit deadline — as in pollUntil above — is what keeps these tests deterministic; a fixed Thread.sleep() guessing how long a rebalance takes is exactly the kind of flaky-test debt that makes a team stop trusting the suite and start skipping it under deadline pressure. 


Where each tool actually belongs 

The right shape for a Kafka test suite isn't "replace embedded tests with Testcontainers everywhere" — that's slower for no benefit on the tests that were never exercising coordination behavior in the first place. Unit tests with a mocked Consumer/Producer cover business logic in isolation, fast, and should be the majority of the suite. @EmbeddedKafka earns its place for exactly what it's good at: confirming Spring wiring and serialization work, quickly, in CI. Testcontainers earns its place specifically for the tests that exist to prove something about Kafka's own coordination behavior — rebalance handling, offset-commit timing, consumer-group recovery — because that's the one category of bug where a lighter-weight substitute isn't just faster, it's testing a different, friendlier system than the one actually running in production.