Topic 7.4
Testing Kafka Applications
In one line
Unit-test processing logic without Kafka, test Kafka Streams topologies with TopologyTestDriver, and run integration tests against a real broker with Testcontainers (preferred) or Spring's @EmbeddedKafka. Test the failure paths: retries, DLT, duplicates and rebalances.
Think of it like this
A fire drill in a real building versus a diagram on paper. Unit tests check the plan; integration tests with a real broker check that people can actually get out.
Key ideas
- 01
Unit tests: keep handlers as plain functions of (record) → effects; test them with plain objects. Most logic doesn't need a broker.
- 02
Kafka Streams:
TopologyTestDriverruns a topology in memory with input and output test topics and fast, deterministic time control (advance wall-clock and stream time to test windows). - 03
Integration: Testcontainers'
KafkaContainer(for exampleapache/kafka-nativeimages for fast startup) runs a real broker per test class; Spring Boot 3.1+@ServiceConnectionwires the bootstrap servers automatically.@EmbeddedKafkaruns an in-JVM broker, faster to start but less realistic. - 04
What to test: the happy path, a poison record reaching the DLT, a transient failure retried, duplicate delivery producing one effect (send the same record twice), consumer restart resuming from the committed offset, and schema compatibility (with a Schema Registry container).
- 05
Use Awaitility for asynchronous assertions instead of sleeps.
Code & diagrams
@SpringBootTest
@Testcontainers
class PaymentListenerIT {
@Container @ServiceConnection
static KafkaContainer kafka = new KafkaContainer(DockerImageName.parse("apache/kafka-native:4.1.0"));
@Autowired KafkaTemplate<String, OrderEvent> template;
@Autowired PaymentRepository payments;
@MockitoBean PaymentProvider provider;
@Test
void duplicateEventChargesOnce() {
OrderEvent evt = TestEvents.orderCreated("order-9");
template.send("commerce.orders", evt.orderId(), evt);
template.send("commerce.orders", evt.orderId(), evt); // redelivery simulation
await().atMost(Duration.ofSeconds(10))
.untilAsserted(() -> assertThat(payments.countByOrderId("order-9")).isEqualTo(1));
verify(provider, atMost(2)).charge(any(), any(), eq("pay-order-9")); // same idempotency key
}
}Explain it without notes
Why is Testcontainers usually preferred over an embedded broker?
Practice
Write an integration test proving a malformed record lands in the DLT and the next valid record is processed.
Trade-offs
- ↔
Container-based tests are slower but catch real integration issues; keep them focused and run unit tests for logic.
Done when you can
I test Kafka code at unit, topology and integration level, including failure paths.