import org.apache.kafka.streams.TestOutputTopic;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.TopologyTestDriver;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.kafka.support.serializer.JsonSerde;
public final static String IN = "TEST-IN";
public final static String OUT = "TEST-OUT";
- @Test
- public void test()
+
+ TopologyTestDriver testDriver;
+ TestInputTopic<Key, Entry> in;
+ TestOutputTopic<String, Ranking> out;
+
+
+ @BeforeEach
+ public void setUp()
{
Topology topology = Top10StreamProcessor.buildTopology(IN, OUT);
JsonSerde<?> valueSerde = new JsonSerde<>();
valueSerde.configure(propertyMap, false);
- TopologyTestDriver testDriver = new TopologyTestDriver(topology, streamProcessorProperties);
+ testDriver = new TopologyTestDriver(topology, streamProcessorProperties);
- TestInputTopic<Key, Entry> in = testDriver.createInputTopic(
+ in = testDriver.createInputTopic(
IN,
(JsonSerializer<Key>)keySerde.serializer(),
(JsonSerializer<Entry>)valueSerde.serializer());
- TestOutputTopic<String, Ranking> out = testDriver.createOutputTopic(
+ out = testDriver.createOutputTopic(
OUT,
(JsonDeserializer<String>)keySerde.deserializer(),
(JsonDeserializer<Ranking>)valueSerde.deserializer());
+ }
+
+
+ @Test
+ public void test()
+ {
Stream
.of(TestData.INPUT_MESSAGES)
.forEach(kv -> in.pipeInput(
});
TestData.assertExpectedMessages(receivedMessages);
+ }
+ @AfterEach
+ public void tearDown()
+ {
testDriver.close();
}
}