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();
         }
     }

Reply via email to