Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

iTechGuides is reader-supported. When you buy through links on our site, we may earn an affiliate commission. As an Amazon Associate I earn from qualifying purchases. Learn more

You can check the logic of a Kafka Streams topology without starting Kafka. Apache Kafka’s TopologyTestDriver, shipped in the kafka-streams-test-utils artifact, runs your topology inside the test process, feeds it records, and lets you read the output and inspect state stores. Apache Kafka’s API documentation describes the class this way: “Best of all, the class works without a real Kafka broker, so the tests execute very quickly with very little overhead.” That speed is the reason to use it. It is also the reason it cannot answer every question about a streaming application, which is what the rest of this article covers.

What the driver can and cannot prove

The driver tests topology logic: how records are transformed, routed, aggregated, and written to output topics, and what ends up in your state stores. It simulates the Kafka consumers and producers the topology would use, and its test helpers convert ordinary Java objects to and from serialized bytes, so your serdes are exercised too.

It does not test how the application behaves against a real cluster. The table below sets the two approaches side by side.

Free tools Windows power users keep installed

One-click scans. No signup required.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Question TopologyTestDriver Broker-backed integration test
Does the topology transform and route records correctly? Yes, the primary use Yes, but more setup per run
Is a running Kafka broker required? No Yes
Speed Described by Apache Kafka as executing “very quickly with very little overhead” Not stated in the cited Apache Kafka sources
Partition behavior Input topics are simulated as single-partitioned Multiple partitions can be exercised against a real cluster
Deployment configuration and broker interaction Not tested Can be tested
State store inspection Available through the driver API Requires querying the running application
Control of wall-clock time Explicit: you advance mocked time Driven by real time

Use the driver for fast, deterministic checks of topology logic. Keep at least a small number of broker-backed tests for anything that depends on partitioning, consumer group behavior, broker configuration, or the deployment environment.

Build a minimal test

1. Add the test dependency

Add kafka-streams-test-utils with the test scope. Apache Kafka’s 3.8 documentation shows it as a test-scoped Maven dependency. Its example version is only an example: set the version to match the kafka-streams version your application already uses.

<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-streams-test-utils</artifactId>
  <version>${kafka.version}</version>
  <scope>test</scope>
</dependency>

2. Build the topology and create the driver

You can test a topology built with the DSL or with the Processor API. With the DSL, construct it with StreamsBuilder and call build(). Then pass the topology and a set of properties to the driver. The properties should include the settings your topology depends on: at minimum the application ID and default key and value serdes, plus a timestamp extractor if your logic relies on one. The bootstrap server value is a placeholder, because the driver does not contact a broker.

import java.util.Properties;
import java.time.Duration;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.TopologyTestDriver;
import org.apache.kafka.streams.TestInputTopic;
import org.apache.kafka.streams.TestOutputTopic;
import org.apache.kafka.streams.state.KeyValueStore;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;

class OrderTotalsTest {

    @Test
    void writesTotalPerOrder() {
        StreamsBuilder builder = new StreamsBuilder();
        // ... define your topology on builder here ...
        Topology topology = builder.build();

        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-totals-test");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        try (TopologyTestDriver driver = new TopologyTestDriver(topology, props)) {
            // steps 3 to 5 go here
        }
    }
}

3. Pipe input and read output

Create a TestInputTopic for each source topic and a TestOutputTopic for each sink. Each helper takes the serializer or deserializer that matches the data. Pipe records into the input, then read the output and assert on it.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
TestInputTopic<String, String> orders = driver.createInputTopic(
    "orders", Serdes.String().serializer(), Serdes.String().serializer());
TestOutputTopic<String, String> totals = driver.createOutputTopic(
    "order-totals", Serdes.String().deserializer(), Serdes.String().deserializer());

orders.pipeInput("order-1", "42");

KeyValue<String, String> result = totals.readKeyValue();
assertEquals("order-1", result.key);
assertEquals("42", result.value);

Read the output in the order your topology writes it. If a test expects a record and none arrives, check the topic names before anything else; a mismatch between the name in the test and the name in the topology is the most common cause of an empty output.

4. Handle time

Time behaves differently depending on the punctuation type your topology uses.

  • Event-time punctuation is triggered by pipeInput. Supply records with timestamps that cover the interval you want to test, such as orders.pipeInput("order-2", "17", timestampMs), and check the result.
  • Wall-clock punctuation does not advance on its own in tests. Advance the driver’s mocked wall clock explicitly, for example driver.advanceWallClockTime(Duration.ofMinutes(5));, and then assert on the output.

5. Inspect state and close the driver

State stores can be read from the test. You can also pre-populate them before input is piped, which is useful for testing how a topology handles existing data. Use the store name defined in the topology.

Rank #4
Metamorphosis: Franz Kafka (Little Clothbound Classics)
  • Metamorphosis: Franz Kafka (Little Clothbound Classics)
KeyValueStore<String, String> store = driver.getKeyValueStore("order-totals-store");
store.put("order-0", "10");   // pre-populate before piping input
// ... pipe input, then check the store ...
assertEquals("42", store.get("order-1"));

Close the driver when the test finishes. The try-with-resources block in the example above does this. Leaving drivers open across tests can leak resources and make failures harder to read.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Settings that do not apply in tests

Apache Kafka’s current API documentation says input processing in the driver is synchronous. Because of that, commit.interval.ms and cache.max.bytes.buffering have no effect. Each input behaves as if it were committed and flushed immediately. A test that depends on caching or commit batching is not exercising those behaviors, so tests that look at intermediate states may not match production.

The partition limit

The most important limit is that input topics are simulated as single-partitioned. Any logic that depends on how records are spread across partitions, such as keying assumptions, ordering across partitions, or repartitioning effects, cannot be confirmed with the driver alone. Write those checks as broker-backed integration tests against a cluster that has multiple partitions.

Troubleshooting checklist

  • No output record: confirm the topic names in the test match the topology exactly.
  • Serialization exception: check that the serdes passed to createInputTopic and createOutputTopic match the data, and that the default serde properties match what the topology expects.
  • Unexpected windowing or time result: confirm the timestamps you supplied and the timestamp extractor configured in the properties.
  • Wall-clock punctuator never fires: call advanceWallClockTime with an interval large enough to reach the punctuation schedule.
  • Store lookup returns null: confirm the store name passed to getKeyValueStore matches the name in the topology.
  • Test passes but production misbehaves: check whether the behavior depends on partitions, commit intervals, caching, or broker configuration, since the driver does not model those.

Versions

The primary API reference is the Apache Kafka 4.3.1 documentation for TopologyTestDriver. The dependency setup example comes from the Apache Kafka 3.8 documentation, which points readers to the latest documentation for current details. Pin the test artifact to the Kafka version your application uses, and check the API documentation for that version when a method signature differs from the examples here.

”

The Bottom Line

Use TopologyTestDriver for fast, deterministic checks of topology logic, and keep broker-backed integration tests for partitioning, deployment, and real broker behavior.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.