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;
@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";
TopologyTestDriver testDriver;
TestInputTopic<Key, Entry> in;
- TestOutputTopic<String, Ranking> out;
+ 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();
out = testDriver.createOutputTopic(
OUT,
- (JsonDeserializer<String>)keySerde.deserializer(),
+ (JsonDeserializer<User>)keySerde.deserializer(),
(JsonDeserializer<Ranking>)valueSerde.deserializer());
}
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 ->
});
TestData.assertExpectedMessages(receivedMessages);
+
+ KeyValueStore<User, Ranking> store = testDriver.getKeyValueStore(STORE_NAME);
+ TestData.assertExpectedState(store);
}
@AfterEach