Command Palette

Search for a command to run...

Hectal
PHASE 7Intermediate ~7 min· topic 4 of 4

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.

0/4 · 0%

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

  1. 01

    Unit tests: keep handlers as plain functions of (record) → effects; test them with plain objects. Most logic doesn't need a broker.

  2. 02

    Kafka Streams: TopologyTestDriver runs a topology in memory with input and output test topics and fast, deterministic time control (advance wall-clock and stream time to test windows).

  3. 03

    Integration: Testcontainers' KafkaContainer (for example apache/kafka-native images for fast startup) runs a real broker per test class; Spring Boot 3.1+ @ServiceConnection wires the bootstrap servers automatically. @EmbeddedKafka runs an in-JVM broker, faster to start but less realistic.

  4. 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).

  5. 05

    Use Awaitility for asynchronous assertions instead of sleeps.

Code & diagrams

PaymentListenerIT.javajava
@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

01

Why is Testcontainers usually preferred over an embedded broker?

Practice

01

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.