import org.apache.kafka.streams.TestOutputTopic;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.TopologyTestDriver;
+import org.apache.kafka.streams.state.KeyValueStore;
+import org.apache.kafka.streams.state.Stores;
+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;
import java.util.Properties;
import java.util.stream.Stream;
-import static de.juplo.kafka.wordcount.top10.TestData.convertToMap;
-import static de.juplo.kafka.wordcount.top10.TestData.parseHeader;
-import static org.springframework.kafka.support.mapping.AbstractJavaTypeMapper.DEFAULT_CLASSID_FIELD_NAME;
-import static org.springframework.kafka.support.mapping.AbstractJavaTypeMapper.KEY_DEFAULT_CLASSID_FIELD_NAME;
+import static de.juplo.kafka.wordcount.top10.Top10ApplicationConfiguration.serializationConfig;
@Slf4j
public class Top10StreamProcessorTopologyTest
{
- public final static String IN = "TEST-IN";
- public final static String OUT = "TEST-OUT";
+ public static final String IN = "TEST-IN";
+ public static final String OUT = "TEST-OUT";
+ public static final String STORE_NAME = "TOPOLOGY-TEST";
- @Test
- public void test()
+
+ TopologyTestDriver testDriver;
+ TestInputTopic<Key, Entry> in;
+ TestOutputTopic<User, Ranking> out;
+
+
+ @BeforeEach
+ public void setUp()
{
- Topology topology = Top10StreamProcessor.buildTopology(IN, OUT);
+ Topology topology = Top10StreamProcessor.buildTopology(
+ IN,
+ OUT,
+ Stores.inMemoryKeyValueStore(STORE_NAME));
- Top10ApplicationConfiguration applicationConfiguriation =
- new Top10ApplicationConfiguration();
- Properties streamProcessorProperties =
- applicationConfiguriation.streamProcessorProperties(new Top10ApplicationProperties());
- Map<String, Object> propertyMap = convertToMap(streamProcessorProperties);
+ Map<String, Object> propertyMap = serializationConfig();
+
+ Properties properties = new Properties();
+ properties.putAll(propertyMap);
JsonSerde<?> keySerde = new JsonSerde<>();
keySerde.configure(propertyMap, true);
JsonSerde<?> valueSerde = new JsonSerde<>();
valueSerde.configure(propertyMap, false);
- TopologyTestDriver testDriver = new TopologyTestDriver(topology, streamProcessorProperties);
+ testDriver = new TopologyTestDriver(topology, properties);
- 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<User>)keySerde.deserializer(),
(JsonDeserializer<Ranking>)valueSerde.deserializer());
+ }
+
+
+ @Test
+ public void test()
+ {
Stream
.of(TestData.INPUT_MESSAGES)
.forEach(kv -> in.pipeInput(
Key.of(kv.key.getUser(), kv.key.getWord()),
Entry.of(kv.value.getWord(), kv.value.getCounter())));
- MultiValueMap<String, Ranking> receivedMessages = new LinkedMultiValueMap<>();
+ MultiValueMap<User, Ranking> receivedMessages = new LinkedMultiValueMap<>();
out
.readRecordsToList()
- .forEach(record ->
- {
- log.debug(
- "OUT: {} -> {}, {}, {}",
- record.key(),
- record.value(),
- parseHeader(record.headers(), KEY_DEFAULT_CLASSID_FIELD_NAME),
- parseHeader(record.headers(), DEFAULT_CLASSID_FIELD_NAME));
- receivedMessages.add(record.key(), record.value());
- });
+ .forEach(record -> receivedMessages.add(record.key(), record.value()));
TestData.assertExpectedMessages(receivedMessages);
+ TestData.assertExpectedNumberOfMessagesForUsers(receivedMessages);
+ TestData.assertExpectedLastMessagesForUsers(receivedMessages);
+
+ KeyValueStore<User, Ranking> store = testDriver.getKeyValueStore(STORE_NAME);
+ TestData.assertExpectedState(store);
+ }
+
+ @AfterEach
+ public void tearDown()
+ {
testDriver.close();
}
}