Testing Topologies

Testing Kafka Streams Topologies

TopologyTestDriver (from kafka-streams-test-utils) runs a topology in-process without a broker: pipe records into input topics, then read output topics and stores, with time under the test's control. streams/test.sh runs this test with the JUnit 48,319 console launcher:

streams/test/booknest/GenreRevenueTest.java: a JUnit 6 test of the topologyJava
package booknest;
import static org.junit.jupiter.api.Assertions.assertEquals;
import java.time.Instant;
import java.util.Properties;
import org.apache.kafka.common.serialization.*;
import org.apache.kafka.streams.*;
import org.junit.jupiter.api.Test;
class GenreRevenueTest {
  @Test void sumsDeliveredOrdersPerGenreAndDay() {
    Properties p = new Properties();
    p.put(StreamsConfig.APPLICATION_ID_CONFIG, "test");
    try (var driver = new TopologyTestDriver(GenreRevenue.build(), p)) {   // no broker
      var s = new StringSerializer();
      var catalog = driver.createInputTopic("booknest.catalog", s, s);
      catalog.pipeInput("1", "{\"genre\":\"Fiction\"}");
      catalog.pipeInput("3", "{\"genre\":\"Cooking\"}");
      var orders = driver.createInputTopic("booknest.orders", s, s);
      for (String status : new String[] {"delivered", "cancelled", "delivered"}) {
        orders.pipeInput("1", "{\"order_ts\":\"2025-03-01T10:00:00Z\",\"status\":\"" + status
            + "\",\"items\":[{\"book_id\":1,\"qty\":2,\"unit_price\":14.99},"
            + "{\"book_id\":3,\"qty\":1,\"unit_price\":24.0}]}");
      }
      var totals = driver.createOutputTopic("booknest.genre-revenue", new StringDeserializer(),
          new DoubleDeserializer()).readKeyValuesToMap();
      assertEquals(59.96, totals.get("Fiction"), 1e-9);       // the cancelled order is skipped
      assertEquals(48.0, totals.get("Cooking"), 1e-9);
      long day = Instant.parse("2025-03-01T00:00:00Z").toEpochMilli();
      assertEquals(59.96, driver.<String, Double>getWindowStore("genre-daily")
          .fetch("Fiction", day), 1e-9);
    }
  }
}
Output
│  └─ GenreRevenueTest ✔
│     └─ sumsDeliveredOrdersPerGenreAndDay() ✔
...
[         1 tests successful      ]

The driver ignores partitioning and rebalances; test those against a real broker (Testcontainers 328,662 ' Kafka 129 module).