Repository: kafka Updated Branches: refs/heads/0.10.2 b0742c202 -> ec3975978
KAFKA-3502: KStreamTestDriver needs to be closed after the test case Found a few recently added unit tests did not close KStreamTestDriver after the test itself is closed; this can cause RocksDB virtual function called if the contained topology has some persistent store since they will be initialized but not closed in time. MINOR fix: found that when closing KStreamTestDriver, we need to first flushing all stores before closing any of them; this is triggered from the `KTableKTableLeftJoin.shouldNotThrowIllegalStateExceptionWhenMultiCacheEvictions`. MINOR fix: in CachingXXXStore, the `name` field is actually used as the cache's namespace, not really the store name or its corresponding topic name. Fixed it by renaming it to `cacheName` and use `this.name()` elsewhere which will call the underlying store's name. Author: Guozhang Wang <[email protected]> Reviewers: Ewen Cheslack-Postava, Eno Thereska, Damian Guy Closes #2432 from guozhangwang/K3502-kstream-builder-test (cherry picked from commit 2ca3e59bb6054cc8ecf6926a74859efe3df000ba) Signed-off-by: Guozhang Wang <[email protected]> Project: http://git-wip-us.apache.org/repos/asf/kafka/repo Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/ec397597 Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/ec397597 Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/ec397597 Branch: refs/heads/0.10.2 Commit: ec3975978582dfec63986b67c38d6e3bb7742f81 Parents: b0742c2 Author: Guozhang Wang <[email protected]> Authored: Wed Jan 25 11:11:21 2017 -0800 Committer: Guozhang Wang <[email protected]> Committed: Wed Jan 25 11:11:32 2017 -0800 ---------------------------------------------------------------------- .../state/internals/CachingKeyValueStore.java | 26 +++++++------- .../state/internals/CachingSessionStore.java | 16 ++++----- .../streams/state/internals/RocksDBStore.java | 1 + .../internals/GlobalKTableJoinsTest.java | 19 +++++++--- .../internals/KGroupedStreamImplTest.java | 37 +++++++++++--------- .../internals/KGroupedTableImplTest.java | 27 +++++++------- .../kstream/internals/KTableAggregateTest.java | 8 ++--- .../internals/KTableKTableLeftJoinTest.java | 5 ++- .../internals/KeyValuePrinterProcessorTest.java | 10 +++--- .../apache/kafka/test/KStreamTestDriver.java | 5 ++- 10 files changed, 87 insertions(+), 67 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java ---------------------------------------------------------------------- diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java index 9a0a976..1e91b47 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingKeyValueStore.java @@ -37,7 +37,7 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, private final Serde<K> keySerde; private final Serde<V> valueSerde; private CacheFlushListener<K, V> flushListener; - private String name; + private String cacheName; private ThreadCache cache; private InternalProcessorContext context; private StateSerdes<K, V> serdes; @@ -73,9 +73,9 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, keySerde == null ? (Serde<K>) context.keySerde() : keySerde, valueSerde == null ? (Serde<V>) context.valueSerde() : valueSerde); - this.name = context.taskId() + "-" + underlying.name(); + this.cacheName = context.taskId() + "-" + underlying.name(); this.cache = this.context.getCache(); - cache.addDirtyEntryFlushListener(name, new ThreadCache.DirtyEntryFlushListener() { + cache.addDirtyEntryFlushListener(cacheName, new ThreadCache.DirtyEntryFlushListener() { @Override public void apply(final List<ThreadCache.DirtyEntry> entries) { for (ThreadCache.DirtyEntry entry : entries) { @@ -108,7 +108,7 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, @Override public synchronized void flush() { - cache.flush(name); + cache.flush(cacheName); underlying.flush(); } @@ -116,7 +116,7 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, public void close() { flush(); underlying.close(); - cache.close(name); + cache.close(cacheName); } @Override @@ -141,12 +141,12 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, private void validateStoreOpen() { if (!isOpen()) { - throw new InvalidStateStoreException("Store " + this.name + " is currently closed"); + throw new InvalidStateStoreException("Store " + this.name() + " is currently closed"); } } private V get(final byte[] rawKey) { - final LRUCacheEntry entry = cache.get(name, rawKey); + final LRUCacheEntry entry = cache.get(cacheName, rawKey); if (entry == null) { final byte[] rawValue = underlying.get(Bytes.wrap(rawKey)); if (rawValue == null) { @@ -155,7 +155,7 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, // only update the cache if this call is on the streamThread // as we don't want other threads to trigger an eviction/flush if (Thread.currentThread().equals(streamThread)) { - cache.put(name, rawKey, new LRUCacheEntry(rawValue)); + cache.put(cacheName, rawKey, new LRUCacheEntry(rawValue)); } return serdes.valueFrom(rawValue); } @@ -173,15 +173,15 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, final byte[] origFrom = serdes.rawKey(from); final byte[] origTo = serdes.rawKey(to); final KeyValueIterator<Bytes, byte[]> storeIterator = underlying.range(Bytes.wrap(origFrom), Bytes.wrap(origTo)); - final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = cache.range(name, origFrom, origTo); + final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = cache.range(cacheName, origFrom, origTo); return new MergedSortedCacheKeyValueStoreIterator<>(cacheIterator, storeIterator, serdes); } @Override public KeyValueIterator<K, V> all() { validateStoreOpen(); - final KeyValueIterator<Bytes, byte[]> storeIterator = new DelegatingPeekingKeyValueIterator<>(name, underlying.all()); - final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = cache.all(name); + final KeyValueIterator<Bytes, byte[]> storeIterator = new DelegatingPeekingKeyValueIterator<>(this.name(), underlying.all()); + final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = cache.all(cacheName); return new MergedSortedCacheKeyValueStoreIterator<>(cacheIterator, storeIterator, serdes); } @@ -199,7 +199,7 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, private synchronized void put(final byte[] rawKey, final V value) { final byte[] rawValue = serdes.rawValue(value); - cache.put(name, rawKey, new LRUCacheEntry(rawValue, true, context.offset(), + cache.put(cacheName, rawKey, new LRUCacheEntry(rawValue, true, context.offset(), context.timestamp(), context.partition(), context.topic())); } @@ -226,7 +226,7 @@ class CachingKeyValueStore<K, V> implements WrappedStateStore, KeyValueStore<K, validateStoreOpen(); final byte[] rawKey = serdes.rawKey(key); final V v = get(rawKey); - cache.delete(name, serdes.rawKey(key)); + cache.delete(cacheName, serdes.rawKey(key)); underlying.delete(Bytes.wrap(rawKey)); return v; } http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingSessionStore.java ---------------------------------------------------------------------- diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingSessionStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingSessionStore.java index fec6609..2cea915 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingSessionStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/CachingSessionStore.java @@ -41,7 +41,7 @@ class CachingSessionStore<K, AGG> extends WrappedStateStore.AbstractWrappedState private Serde<K> keySerde; private final Serde<AGG> aggSerde; private InternalProcessorContext context; - private String name; + private String cacheName; private StateSerdes<K, AGG> serdes; private ThreadCache cache; private CacheFlushListener<Windowed<K>, AGG> flushListener; @@ -60,8 +60,8 @@ class CachingSessionStore<K, AGG> extends WrappedStateStore.AbstractWrappedState final long earliestSessionEndTime, final long latestSessionStartTime) { validateStoreOpen(); - final Bytes binarySessionId = Bytes.wrap(keySerde.serializer().serialize(name, key)); - final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = cache.range(name, + final Bytes binarySessionId = Bytes.wrap(keySerde.serializer().serialize(this.name(), key)); + final ThreadCache.MemoryLRUCacheBytesIterator cacheIterator = cache.range(cacheName, keySchema.lowerRange(binarySessionId, earliestSessionEndTime).get(), keySchema.upperRange(binarySessionId, latestSessionStartTime).get()); @@ -84,7 +84,7 @@ class CachingSessionStore<K, AGG> extends WrappedStateStore.AbstractWrappedState final Bytes binaryKey = SessionKeySerde.toBinary(key, keySerde.serializer()); final LRUCacheEntry entry = new LRUCacheEntry(serdes.rawValue(value), true, context.offset(), key.window().end(), context.partition(), context.topic()); - cache.put(name, binaryKey.get(), entry); + cache.put(cacheName, binaryKey.get(), entry); } @Override @@ -107,9 +107,9 @@ class CachingSessionStore<K, AGG> extends WrappedStateStore.AbstractWrappedState aggSerde == null ? (Serde<AGG>) context.valueSerde() : aggSerde); - this.name = context.taskId() + "-" + bytesStore.name(); + this.cacheName = context.taskId() + "-" + bytesStore.name(); this.cache = this.context.getCache(); - cache.addDirtyEntryFlushListener(name, new ThreadCache.DirtyEntryFlushListener() { + cache.addDirtyEntryFlushListener(cacheName, new ThreadCache.DirtyEntryFlushListener() { @Override public void apply(final List<ThreadCache.DirtyEntry> entries) { for (ThreadCache.DirtyEntry entry : entries) { @@ -150,14 +150,14 @@ class CachingSessionStore<K, AGG> extends WrappedStateStore.AbstractWrappedState public void flush() { - cache.flush(name); + cache.flush(cacheName); bytesStore.flush(); } public void close() { flush(); bytesStore.close(); - cache.close(name); + cache.close(cacheName); } public void setFlushListener(CacheFlushListener<Windowed<K>, AGG> flushListener) { http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java ---------------------------------------------------------------------- diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java index 55c1bb8..5c83ddf 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java @@ -352,6 +352,7 @@ public class RocksDBStore<K, V> implements KeyValueStore<K, V> { if (!open) { return; } + open = false; closeOpenIterators(); options.close(); http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/streams/kstream/internals/GlobalKTableJoinsTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/GlobalKTableJoinsTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/GlobalKTableJoinsTest.java index bbc9741..c41d953 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/GlobalKTableJoinsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/GlobalKTableJoinsTest.java @@ -25,6 +25,7 @@ import org.apache.kafka.streams.kstream.KeyValueMapper; import org.apache.kafka.test.KStreamTestDriver; import org.apache.kafka.test.MockValueJoiner; import org.apache.kafka.test.TestUtils; +import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -38,14 +39,15 @@ import static org.junit.Assert.assertEquals; public class GlobalKTableJoinsTest { private final KStreamBuilder builder = new KStreamBuilder(); - private GlobalKTable<String, String> global; - private File stateDir; private final Map<String, String> results = new HashMap<>(); + private final String streamTopic = "stream"; + private final String globalTopic = "global"; + private File stateDir; + private GlobalKTable<String, String> global; private KStream<String, String> stream; private KeyValueMapper<String, String, String> keyValueMapper; private ForeachAction<String, String> action; - private final String streamTopic = "stream"; - private final String globalTopic = "global"; + private KStreamTestDriver driver = null; @Before public void setUp() throws Exception { @@ -64,7 +66,14 @@ public class GlobalKTableJoinsTest { results.put(key, value); } }; + } + @After + public void cleanup() { + if (driver != null) { + driver.close(); + } + driver = null; } @Test @@ -94,7 +103,7 @@ public class GlobalKTableJoinsTest { } private void verifyJoin(final Map<String, String> expected, final String joinInput) { - final KStreamTestDriver driver = new KStreamTestDriver(builder, stateDir); + driver = new KStreamTestDriver(builder, stateDir); driver.setTime(0L); // write some data to the global table driver.process(globalTopic, "a", "A"); http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImplTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImplTest.java index b6d8a97..3fd287d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedStreamImplTest.java @@ -32,12 +32,14 @@ import org.apache.kafka.streams.kstream.TimeWindows; import org.apache.kafka.streams.processor.StateStoreSupplier; import org.apache.kafka.streams.kstream.Windowed; import org.apache.kafka.streams.kstream.Windows; +import org.apache.kafka.streams.state.KeyValueStore; import org.apache.kafka.streams.state.SessionStore; import org.apache.kafka.test.KStreamTestDriver; import org.apache.kafka.test.MockAggregator; import org.apache.kafka.test.MockInitializer; import org.apache.kafka.test.MockReducer; import org.apache.kafka.test.TestUtils; +import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -56,6 +58,7 @@ public class KGroupedStreamImplTest { private static final String TOPIC = "topic"; private final KStreamBuilder builder = new KStreamBuilder(); private KGroupedStream<String, String> groupedStream; + private KStreamTestDriver driver = null; @Before public void before() { @@ -63,6 +66,14 @@ public class KGroupedStreamImplTest { groupedStream = stream.groupByKey(Serdes.String(), Serdes.String()); } + @After + public void cleanup() { + if (driver != null) { + driver.close(); + } + driver = null; + } + @Test(expected = NullPointerException.class) public void shouldNotHaveNullReducerOnReduce() throws Exception { groupedStream.reduce(null, "store"); @@ -70,14 +81,12 @@ public class KGroupedStreamImplTest { @Test(expected = NullPointerException.class) public void shouldNotHaveNullStoreNameOnReduce() throws Exception { - String storeName = null; - groupedStream.reduce(MockReducer.STRING_ADDER, storeName); + groupedStream.reduce(MockReducer.STRING_ADDER, (String) null); } @Test(expected = NullPointerException.class) public void shouldNotHaveNullStoreSupplierOnReduce() throws Exception { - StateStoreSupplier storeSupplier = null; - groupedStream.reduce(MockReducer.STRING_ADDER, storeSupplier); + groupedStream.reduce(MockReducer.STRING_ADDER, (StateStoreSupplier<KeyValueStore>) null); } @Test(expected = NullPointerException.class) @@ -92,8 +101,7 @@ public class KGroupedStreamImplTest { @Test(expected = NullPointerException.class) public void shouldNotHaveNullStoreNameWithWindowedReduce() throws Exception { - String storeName = null; - groupedStream.reduce(MockReducer.STRING_ADDER, TimeWindows.of(10), storeName); + groupedStream.reduce(MockReducer.STRING_ADDER, TimeWindows.of(10), (String) null); } @Test(expected = NullPointerException.class) @@ -108,8 +116,7 @@ public class KGroupedStreamImplTest { @Test(expected = NullPointerException.class) public void shouldNotHaveNullStoreNameOnAggregate() throws Exception { - String storeName = null; - groupedStream.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, Serdes.String(), storeName); + groupedStream.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, Serdes.String(), null); } @Test(expected = NullPointerException.class) @@ -129,14 +136,12 @@ public class KGroupedStreamImplTest { @Test(expected = NullPointerException.class) public void shouldNotHaveNullStoreNameOnWindowedAggregate() throws Exception { - String storeName = null; - groupedStream.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, TimeWindows.of(10), Serdes.String(), storeName); + groupedStream.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, TimeWindows.of(10), Serdes.String(), null); } @Test(expected = NullPointerException.class) public void shouldNotHaveNullStoreSupplierOnWindowedAggregate() throws Exception { - StateStoreSupplier storeSupplier = null; - groupedStream.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, TimeWindows.of(10), storeSupplier); + groupedStream.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, TimeWindows.of(10), null); } @Test @@ -165,7 +170,7 @@ public class KGroupedStreamImplTest { } }); - final KStreamTestDriver driver = new KStreamTestDriver(builder, TestUtils.tempDirectory()); + driver = new KStreamTestDriver(builder, TestUtils.tempDirectory()); driver.setTime(10); driver.process(TOPIC, "1", "1"); driver.setTime(15); @@ -194,7 +199,7 @@ public class KGroupedStreamImplTest { results.put(key, value); } }); - final KStreamTestDriver driver = new KStreamTestDriver(builder, TestUtils.tempDirectory()); + driver = new KStreamTestDriver(builder, TestUtils.tempDirectory()); driver.setTime(10); driver.process(TOPIC, "1", "1"); driver.setTime(15); @@ -230,7 +235,7 @@ public class KGroupedStreamImplTest { results.put(key, value); } }); - final KStreamTestDriver driver = new KStreamTestDriver(builder, TestUtils.tempDirectory()); + driver = new KStreamTestDriver(builder, TestUtils.tempDirectory()); driver.setTime(10); driver.process(TOPIC, "1", "A"); driver.setTime(15); @@ -357,7 +362,7 @@ public class KGroupedStreamImplTest { } }); - final KStreamTestDriver driver = new KStreamTestDriver(builder, TestUtils.tempDirectory(), 0); + driver = new KStreamTestDriver(builder, TestUtils.tempDirectory(), 0); driver.setTime(0); driver.process(TOPIC, "1", "A"); driver.process(TOPIC, "2", "B"); http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java index 5ed61e1..5fded02 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KGroupedTableImplTest.java @@ -24,13 +24,13 @@ import org.apache.kafka.streams.kstream.KGroupedTable; import org.apache.kafka.streams.kstream.KStreamBuilder; import org.apache.kafka.streams.kstream.KTable; import org.apache.kafka.streams.kstream.KeyValueMapper; -import org.apache.kafka.streams.processor.StateStoreSupplier; import org.apache.kafka.test.KStreamTestDriver; import org.apache.kafka.test.MockAggregator; import org.apache.kafka.test.MockInitializer; import org.apache.kafka.test.MockKeyValueMapper; import org.apache.kafka.test.MockReducer; import org.apache.kafka.test.TestUtils; +import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -42,19 +42,27 @@ import static org.junit.Assert.assertEquals; public class KGroupedTableImplTest { + private final KStreamBuilder builder = new KStreamBuilder(); private KGroupedTable<String, String> groupedTable; + private KStreamTestDriver driver = null; @Before public void before() { - final KStreamBuilder builder = new KStreamBuilder(); groupedTable = builder.table(Serdes.String(), Serdes.String(), "blah", "blah") .groupBy(MockKeyValueMapper.<String, String>SelectValueKeyValueMapper()); } + @After + public void cleanup() { + if (driver != null) { + driver.close(); + } + driver = null; + } + @Test(expected = NullPointerException.class) public void shouldNotAllowNullStoreNameOnAggregate() throws Exception { - String storeName = null; - groupedTable.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, MockAggregator.TOSTRING_REMOVER, storeName); + groupedTable.aggregate(MockInitializer.STRING_INIT, MockAggregator.TOSTRING_ADDER, MockAggregator.TOSTRING_REMOVER, (String) null); } @Test(expected = NullPointerException.class) @@ -84,19 +92,16 @@ public class KGroupedTableImplTest { @Test(expected = NullPointerException.class) public void shouldNotAllowNullStoreNameOnReduce() throws Exception { - String storeName = null; - groupedTable.reduce(MockReducer.STRING_ADDER, MockReducer.STRING_REMOVER, storeName); + groupedTable.reduce(MockReducer.STRING_ADDER, MockReducer.STRING_REMOVER, (String) null); } @Test(expected = NullPointerException.class) public void shouldNotAllowNullStoreSupplierOnReduce() throws Exception { - StateStoreSupplier storeName = null; - groupedTable.reduce(MockReducer.STRING_ADDER, MockReducer.STRING_REMOVER, storeName); + groupedTable.reduce(MockReducer.STRING_ADDER, MockReducer.STRING_REMOVER, (String) null); } @Test public void shouldReduce() throws Exception { - final KStreamBuilder builder = new KStreamBuilder(); final String topic = "input"; final KeyValueMapper<String, Number, KeyValue<String, Integer>> intProjection = new KeyValueMapper<String, Number, KeyValue<String, Integer>>() { @@ -118,8 +123,7 @@ public class KGroupedTableImplTest { } }); - - final KStreamTestDriver driver = new KStreamTestDriver(builder, TestUtils.tempDirectory(), Serdes.String(), Serdes.Integer()); + driver = new KStreamTestDriver(builder, TestUtils.tempDirectory(), Serdes.String(), Serdes.Integer()); driver.setTime(10L); driver.process(topic, "A", 1.1); driver.process(topic, "B", 2.2); @@ -136,6 +140,5 @@ public class KGroupedTableImplTest { assertEquals(Integer.valueOf(5), results.get("A")); assertEquals(Integer.valueOf(6), results.get("B")); - } } http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableAggregateTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableAggregateTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableAggregateTest.java index a43cab6..39baa4e 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableAggregateTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableAggregateTest.java @@ -220,7 +220,7 @@ public class KTableAggregateTest { .toStream() .process(proc); - final KStreamTestDriver driver = new KStreamTestDriver(builder, stateDir); + driver = new KStreamTestDriver(builder, stateDir); driver.process(input, "A", "green"); driver.flushState(); @@ -256,7 +256,7 @@ public class KTableAggregateTest { .toStream() .process(proc); - final KStreamTestDriver driver = new KStreamTestDriver(builder, stateDir); + driver = new KStreamTestDriver(builder, stateDir); driver.process(input, "A", "green"); driver.process(input, "B", "green"); @@ -309,7 +309,7 @@ public class KTableAggregateTest { .toStream() .process(proc); - final KStreamTestDriver driver = new KStreamTestDriver(builder, stateDir); + driver = new KStreamTestDriver(builder, stateDir); driver.process(input, "11", "A"); driver.flushState(); @@ -378,7 +378,7 @@ public class KTableAggregateTest { } }); - final KStreamTestDriver driver = new KStreamTestDriver(builder, stateDir, 111); + driver = new KStreamTestDriver(builder, stateDir, 111); driver.process(reduceTopic, "1", new Change<>(1L, null)); driver.process("tableOne", "2", "2"); // this should trigger eviction on the reducer-store topic http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableLeftJoinTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableLeftJoinTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableLeftJoinTest.java index cbf1da4..9790044 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableLeftJoinTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KTableKTableLeftJoinTest.java @@ -367,7 +367,6 @@ public class KTableKTableLeftJoinTest { }; final KTable<Long, String> seven = one.mapValues(mapper); - final KTable<Long, String> eight = six.leftJoin(seven, MockValueJoiner.TOSTRING_JOINER); aggTable.leftJoin(one, MockValueJoiner.TOSTRING_JOINER) @@ -378,7 +377,7 @@ public class KTableKTableLeftJoinTest { .leftJoin(eight, MockValueJoiner.TOSTRING_JOINER) .mapValues(mapper); - final KStreamTestDriver driver = new KStreamTestDriver(builder, stateDir, 250); + driver = new KStreamTestDriver(builder, stateDir, 250); final String[] values = {"a", "AA", "BBB", "CCCC", "DD", "EEEEEEEE", "F", "GGGGGGGGGGGGGGG", "HHH", "IIIIIIIIII", "J", "KK", "LLLL", "MMMMMMMMMMMMMMMMMMMMMM", "NNNNN", "O", "P", "QQQQQ", "R", "SSSS", @@ -387,7 +386,7 @@ public class KTableKTableLeftJoinTest { final Random random = new Random(); for (int i = 0; i < 1000; i++) { for (String input : inputs) { - final Long key = Long.valueOf(random.nextInt(1000)); + final Long key = (long) random.nextInt(1000); final String value = values[random.nextInt(values.length)]; driver.process(input, key, value); } http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KeyValuePrinterProcessorTest.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KeyValuePrinterProcessorTest.java b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KeyValuePrinterProcessorTest.java index d1bff05..4ab1c35 100644 --- a/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KeyValuePrinterProcessorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/kstream/internals/KeyValuePrinterProcessorTest.java @@ -38,11 +38,11 @@ import static org.junit.Assert.assertEquals; public class KeyValuePrinterProcessorTest { - private String topicName = "topic"; - private Serde<String> stringSerde = Serdes.String(); - private ByteArrayOutputStream baos = new ByteArrayOutputStream(); - private KStreamBuilder builder = new KStreamBuilder(); - private PrintStream printStream = new PrintStream(baos); + private final String topicName = "topic"; + private final Serde<String> stringSerde = Serdes.String(); + private final ByteArrayOutputStream baos = new ByteArrayOutputStream(); + private final KStreamBuilder builder = new KStreamBuilder(); + private final PrintStream printStream = new PrintStream(baos); private KStreamTestDriver driver = null; http://git-wip-us.apache.org/repos/asf/kafka/blob/ec397597/streams/src/test/java/org/apache/kafka/test/KStreamTestDriver.java ---------------------------------------------------------------------- diff --git a/streams/src/test/java/org/apache/kafka/test/KStreamTestDriver.java b/streams/src/test/java/org/apache/kafka/test/KStreamTestDriver.java index 011532b..dfa2987 100644 --- a/streams/src/test/java/org/apache/kafka/test/KStreamTestDriver.java +++ b/streams/src/test/java/org/apache/kafka/test/KStreamTestDriver.java @@ -196,8 +196,11 @@ public class KStreamTestDriver { } private void closeState() { + // we need to first flush all stores before trying to close any one + // of them since the flushing could cause eviction and hence tries to access other stores + flushState(); + for (StateStore stateStore : context.allStateStores().values()) { - stateStore.flush(); stateStore.close(); } }
