From 70f7067babdf1c06d805450c947bf777553a53b2 Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Sat, 8 Dec 2018 20:32:29 +0300 Subject: [PATCH 1/7] KAFKA-6970: All standard state stores from ProcessContext are guarded with read only wrapper --- .../internals/ProcessorContextImpl.java | 22 ++++++++++++------- 1 file changed, 14 insertions(+), 8 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java index c79ec35328af8..9408a11b51de3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java @@ -75,20 +75,26 @@ public RecordCollector recordCollector() { @SuppressWarnings("unchecked") @Override public StateStore getStateStore(final String name) { + StateStore stateStore = getStateStore0(name); + + if (stateStore instanceof KeyValueStore) { + return new KeyValueStoreReadOnlyDecorator((KeyValueStore) stateStore); + } else if (stateStore instanceof WindowStore) { + return new WindowStoreReadOnlyDecorator((WindowStore) stateStore); + } else if (stateStore instanceof SessionStore) { + return new SessionStoreReadOnlyDecorator((SessionStore) stateStore); + } + + return stateStore; + } + + private StateStore getStateStore0(final String name) { if (currentNode() == null) { throw new StreamsException("Accessing from an unknown node"); } final StateStore global = stateManager.getGlobalStore(name); if (global != null) { - if (global instanceof KeyValueStore) { - return new KeyValueStoreReadOnlyDecorator((KeyValueStore) global); - } else if (global instanceof WindowStore) { - return new WindowStoreReadOnlyDecorator((WindowStore) global); - } else if (global instanceof SessionStore) { - return new SessionStoreReadOnlyDecorator((SessionStore) global); - } - return global; } From 6ab1684576b1b9df8b3c1fa181fa85a9a280da17 Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Sun, 9 Dec 2018 10:22:37 +0300 Subject: [PATCH 2/7] KAFKA-6970: Reworked to guard local state store with ReadWriteDecorator. --- .../internals/ProcessorContextImpl.java | 200 ++++++++++++++++-- .../internals/ProcessorContextImplTest.java | 200 ++++++++++++++++-- 2 files changed, 365 insertions(+), 35 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java index 9408a11b51de3..a862dc26fe65a 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.processor.internals; +import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.internals.ApiUtils; @@ -75,26 +76,20 @@ public RecordCollector recordCollector() { @SuppressWarnings("unchecked") @Override public StateStore getStateStore(final String name) { - StateStore stateStore = getStateStore0(name); - - if (stateStore instanceof KeyValueStore) { - return new KeyValueStoreReadOnlyDecorator((KeyValueStore) stateStore); - } else if (stateStore instanceof WindowStore) { - return new WindowStoreReadOnlyDecorator((WindowStore) stateStore); - } else if (stateStore instanceof SessionStore) { - return new SessionStoreReadOnlyDecorator((SessionStore) stateStore); - } - - return stateStore; - } - - private StateStore getStateStore0(final String name) { if (currentNode() == null) { throw new StreamsException("Accessing from an unknown node"); } final StateStore global = stateManager.getGlobalStore(name); if (global != null) { + if (global instanceof KeyValueStore) { + return new KeyValueStoreReadOnlyDecorator((KeyValueStore) global); + } else if (global instanceof WindowStore) { + return new WindowStoreReadOnlyDecorator((WindowStore) global); + } else if (global instanceof SessionStore) { + return new SessionStoreReadOnlyDecorator((SessionStore) global); + } + return global; } @@ -108,7 +103,16 @@ private StateStore getStateStore0(final String name) { "please file a bug report at https://issues.apache.org/jira/projects/KAFKA."); } - return stateManager.getStore(name); + final StateStore store = stateManager.getStore(name); + if (store instanceof KeyValueStore) { + return new KeyValueStoreReadWriteDecorator((KeyValueStore) store); + } else if (store instanceof WindowStore) { + return new WindowStoreReadWriteDecorator((WindowStore) store); + } else if (store instanceof SessionStore) { + return new SessionStoreReadWriteDecorator((SessionStore) store); + } + + return store; } @SuppressWarnings("unchecked") @@ -367,4 +371,170 @@ public KeyValueIterator, AGG> fetch(final K from, final K to) { return underlying.fetch(from, to); } } + + private abstract static class StateStoreReadWriteDecorator implements StateStore { + static final String ERROR_MESSAGE = "This method may only be called by Kafka Streams"; + + final T underlying; + + StateStoreReadWriteDecorator(final T underlying) { + this.underlying = underlying; + } + + @Override + public void init(final ProcessorContext context, final StateStore root) { + throw new UnsupportedOperationException(ERROR_MESSAGE); + } + + @Override + public void close() { + throw new UnsupportedOperationException(ERROR_MESSAGE); + } + + @Override + public String name() { + return underlying.name(); + } + + @Override + public void flush() { + underlying.flush(); + } + + @Override + public boolean persistent() { + return underlying.persistent(); + } + + @Override + public boolean isOpen() { + return underlying.isOpen(); + } + } + + private static class KeyValueStoreReadWriteDecorator extends StateStoreReadWriteDecorator> implements KeyValueStore { + KeyValueStoreReadWriteDecorator(final KeyValueStore underlying) { + super(underlying); + } + + @Override + public V get(final K key) { + return underlying.get(key); + } + + @Override + public KeyValueIterator range(final K from, final K to) { + return underlying.range(from, to); + } + + @Override + public KeyValueIterator all() { + return underlying.all(); + } + + @Override + public long approximateNumEntries() { + return underlying.approximateNumEntries(); + } + + @Override + public void put(final K key, final V value) { + underlying.put(key, value); + } + + @Override + public V putIfAbsent(final K key, final V value) { + return underlying.putIfAbsent(key, value); + } + + @Override + public void putAll(final List> entries) { + underlying.putAll(entries); + } + + @Override + public V delete(final K key) { + return underlying.delete(key); + } + } + + private static class WindowStoreReadWriteDecorator extends StateStoreReadWriteDecorator> implements WindowStore { + WindowStoreReadWriteDecorator(final WindowStore underlying) { + super(underlying); + } + + @Override + public void put(final K key, final V value) { + underlying.put(key, value); + } + + @Override + public void put(final K key, final V value, final long windowStartTimestamp) { + underlying.put(key, value, windowStartTimestamp); + } + + @Override + public V fetch(final K key, final long time) { + return underlying.fetch(key, time); + } + + @Deprecated + @Override + public WindowStoreIterator fetch(final K key, final long timeFrom, final long timeTo) { + return underlying.fetch(key, timeFrom, timeTo); + } + + @Deprecated + @Override + public KeyValueIterator, V> fetch(final K from, final K to, final long timeFrom, final long timeTo) { + return underlying.fetch(from, to, timeFrom, timeTo); + } + + @Override + public KeyValueIterator, V> all() { + return underlying.all(); + } + + @Deprecated + @Override + public KeyValueIterator, V> fetchAll(final long timeFrom, final long timeTo) { + return underlying.fetchAll(timeFrom, timeTo); + } + } + + private static class SessionStoreReadWriteDecorator extends StateStoreReadWriteDecorator> implements SessionStore { + SessionStoreReadWriteDecorator(final SessionStore underlying) { + super(underlying); + } + + @Override + public KeyValueIterator, AGG> findSessions(final K key, final long earliestSessionEndTime, final long latestSessionStartTime) { + return underlying.findSessions(key, earliestSessionEndTime, latestSessionStartTime); + } + + @Override + public KeyValueIterator, AGG> findSessions(final K keyFrom, final K keyTo, final long earliestSessionEndTime, final long latestSessionStartTime) { + return underlying.findSessions(keyFrom, keyTo, earliestSessionEndTime, latestSessionStartTime); + } + + @Override + public void remove(final Windowed sessionKey) { + underlying.remove(sessionKey); + } + + @Override + public void put(final Windowed sessionKey, final AGG aggregate) { + underlying.put(sessionKey, aggregate); + } + + @Override + public KeyValueIterator, AGG> fetch(final K key) { + return underlying.fetch(key); + } + + @Override + public KeyValueIterator, AGG> fetch(final K from, final K to) { + return underlying.fetch(from, to); + } + } } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java index fa5f597aa8999..bb6456ba9e24c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java @@ -18,6 +18,7 @@ import java.util.ArrayList; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.function.Consumer; @@ -38,8 +39,9 @@ import org.junit.Before; import org.junit.Test; -import static java.util.Collections.emptySet; +import static java.util.Arrays.asList; import static org.easymock.EasyMock.anyLong; +import static org.easymock.EasyMock.anyObject; import static org.easymock.EasyMock.anyString; import static org.easymock.EasyMock.expect; import static org.easymock.EasyMock.expectLastCall; @@ -58,6 +60,14 @@ public class ProcessorContextImplTest { private boolean initExecuted; private boolean closeExecuted; + private boolean flushExecuted; + private boolean putExecuted; + private boolean putIfAbsentExecuted; + private boolean putAllExecuted; + private boolean deleteExecuted; + private boolean removeExecuted; + private boolean put3argExecuted; + private KeyValueIterator rangeIter; private KeyValueIterator allIter; @@ -66,6 +76,16 @@ public class ProcessorContextImplTest { @Before public void setup() { + initExecuted = false; + closeExecuted = false; + flushExecuted = false; + putExecuted = false; + putIfAbsentExecuted = false; + putAllExecuted = false; + deleteExecuted = false; + removeExecuted = false; + put3argExecuted = false; + rangeIter = mock(KeyValueIterator.class); allIter = mock(KeyValueIterator.class); windowStoreIter = mock(WindowStoreIterator.class); @@ -82,9 +102,14 @@ public void setup() { final ProcessorStateManager stateManager = mock(ProcessorStateManager.class); - expect(stateManager.getGlobalStore("KeyValueStore")).andReturn(keyValueStoreMock()); - expect(stateManager.getGlobalStore("WindowStore")).andReturn(windowStoreMock()); - expect(stateManager.getGlobalStore("SessionStore")).andReturn(sessionStoreMock()); + expect(stateManager.getGlobalStore("GlobalKeyValueStore")).andReturn(keyValueStoreMock()); + expect(stateManager.getGlobalStore("GlobalWindowStore")).andReturn(windowStoreMock()); + expect(stateManager.getGlobalStore("GlobalSessionStore")).andReturn(sessionStoreMock()); + expect(stateManager.getGlobalStore(anyString())).andReturn(null); + + expect(stateManager.getStore("LocalKeyValueStore")).andReturn(keyValueStoreMock()); + expect(stateManager.getStore("LocalWindowStore")).andReturn(windowStoreMock()); + expect(stateManager.getStore("LocalSessionStore")).andReturn(sessionStoreMock()); replay(stateManager); @@ -98,12 +123,15 @@ public void setup() { mock(ThreadCache.class) ); - context.setCurrentNode(new ProcessorNode("fake", null, emptySet())); + context.setCurrentNode(new ProcessorNode("fake", null, + new HashSet<>(asList("LocalKeyValueStore", "LocalWindowStore", "LocalSessionStore")))); } @Test - public void testKeyValueStore() { - doTest("KeyValueStore", (Consumer>) store -> { + public void testGlobalKeyValueStore() { + doTest("GlobalKeyValueStore", (Consumer>) store -> { + checkGlobalStateStoreMethods(store); + checkThrowsUnsupportedOperation(() -> store.put("1", 1L), "put"); checkThrowsUnsupportedOperation(() -> store.putIfAbsent("1", 1L), "putIfAbsent"); checkThrowsUnsupportedOperation(() -> store.putAll(Collections.emptyList()), "putAll"); @@ -117,8 +145,10 @@ public void testKeyValueStore() { } @Test - public void testWindowStore() { - doTest("WindowStore", (Consumer>) store -> { + public void testGlobalWindowStore() { + doTest("GlobalWindowStore", (Consumer>) store -> { + checkGlobalStateStoreMethods(store); + checkThrowsUnsupportedOperation(() -> store.put("1", 1L, 1L), "put"); checkThrowsUnsupportedOperation(() -> store.put("1", 1L), "put"); @@ -131,8 +161,10 @@ public void testWindowStore() { } @Test - public void testSessionStore() { - doTest("SessionStore", (Consumer>) store -> { + public void testGlobalSessionStore() { + doTest("GlobalSessionStore", (Consumer>) store -> { + checkGlobalStateStoreMethods(store); + checkThrowsUnsupportedOperation(() -> store.remove(null), "remove"); checkThrowsUnsupportedOperation(() -> store.put(null, null), "put"); @@ -143,6 +175,68 @@ public void testSessionStore() { }); } + @Test + public void testLocalKeyValueStore() { + doTest("LocalKeyValueStore", (Consumer>) store -> { + checkLocalStateStoreMethods(store); + + store.put("1", 1L); + assertTrue(putExecuted); + + store.putIfAbsent("1", 1L); + assertTrue(putIfAbsentExecuted); + + store.putAll(Collections.emptyList()); + assertTrue(putAllExecuted); + + store.delete("1"); + assertTrue(deleteExecuted); + + assertEquals((Long) VAL, store.get(KEY)); + assertEquals(rangeIter, store.range("one", "two")); + assertEquals(allIter, store.all()); + assertEquals(VAL, store.approximateNumEntries()); + }); + } + + @Test + public void testLocalWindowStore() { + doTest("LocalWindowStore", (Consumer>) store -> { + checkLocalStateStoreMethods(store); + + store.put("1", 1L); + assertTrue(putExecuted); + + store.put("1", 1L, 1L); + assertTrue(put3argExecuted); + + assertEquals(iters.get(0), store.fetchAll(0L, 0L)); + assertEquals(windowStoreIter, store.fetch(KEY, 0L, 1L)); + assertEquals(iters.get(1), store.fetch(KEY, KEY, 0L, 1L)); + assertEquals((Long) VAL, store.fetch(KEY, 1L)); + assertEquals(iters.get(2), store.all()); + }); + } + + @Test + public void testLocalSessionStore() { + doTest("LocalSessionStore", (Consumer>) store -> { + checkLocalStateStoreMethods(store); + + store.remove(null); + assertTrue(removeExecuted); + + store.put(null, null); + assertTrue(putExecuted); + + assertEquals(iters.get(3), store.findSessions(KEY, 1L, 2L)); + assertEquals(iters.get(4), store.findSessions(KEY, KEY, 1L, 2L)); + assertEquals(iters.get(5), store.fetch(KEY)); + assertEquals(iters.get(6), store.fetch(KEY, KEY)); + }); + } + + @SuppressWarnings("unchecked") private KeyValueStore keyValueStoreMock() { final KeyValueStore keyValueStoreMock = mock(KeyValueStore.class); @@ -154,6 +248,31 @@ private KeyValueStore keyValueStoreMock() { expect(keyValueStoreMock.range("one", "two")).andReturn(rangeIter); expect(keyValueStoreMock.all()).andReturn(allIter); + + keyValueStoreMock.put(anyString(), anyLong()); + expectLastCall().andAnswer(() -> { + putExecuted = true; + return null; + }); + + keyValueStoreMock.putIfAbsent(anyString(), anyLong()); + expectLastCall().andAnswer(() -> { + putIfAbsentExecuted = true; + return null; + }); + + keyValueStoreMock.putAll(anyObject(List.class)); + expectLastCall().andAnswer(() -> { + putAllExecuted = true; + return null; + }); + + keyValueStoreMock.delete(anyString()); + expectLastCall().andAnswer(() -> { + deleteExecuted = true; + return null; + }); + replay(keyValueStoreMock); return keyValueStoreMock; @@ -170,11 +289,24 @@ private WindowStore windowStoreMock() { expect(windowStore.fetch(anyString(), anyLong())).andReturn(VAL); expect(windowStore.all()).andReturn(iters.get(2)); + windowStore.put(anyString(), anyLong()); + expectLastCall().andAnswer(() -> { + putExecuted = true; + return null; + }); + + windowStore.put(anyString(), anyLong(), anyLong()); + expectLastCall().andAnswer(() -> { + put3argExecuted = true; + return null; + }); + replay(windowStore); return windowStore; } + @SuppressWarnings("unchecked") private SessionStore sessionStoreMock() { final SessionStore sessionStore = mock(SessionStore.class); @@ -185,27 +317,45 @@ private SessionStore sessionStoreMock() { expect(sessionStore.fetch(anyString())).andReturn(iters.get(5)); expect(sessionStore.fetch(anyString(), anyString())).andReturn(iters.get(6)); + sessionStore.put(anyObject(Windowed.class), anyLong()); + expectLastCall().andAnswer(() -> { + putExecuted = true; + return null; + }); + + sessionStore.remove(anyObject(Windowed.class)); + expectLastCall().andAnswer(() -> { + removeExecuted = true; + return null; + }); + replay(sessionStore); return sessionStore; } - private void initStateStoreMock(final StateStore windowStore) { - expect(windowStore.name()).andReturn(STORE_NAME); - expect(windowStore.persistent()).andReturn(true); - expect(windowStore.isOpen()).andReturn(true); + private void initStateStoreMock(final StateStore stateStore) { + expect(stateStore.name()).andReturn(STORE_NAME); + expect(stateStore.persistent()).andReturn(true); + expect(stateStore.isOpen()).andReturn(true); - windowStore.init(null, null); + stateStore.init(null, null); expectLastCall().andAnswer(() -> { initExecuted = true; return null; }); - windowStore.close(); + stateStore.close(); expectLastCall().andAnswer(() -> { closeExecuted = true; return null; }); + + stateStore.flush(); + expectLastCall().andAnswer(() -> { + flushExecuted = true; + return null; + }); } private void doTest(final String name, final Consumer checker) { @@ -215,8 +365,6 @@ private void doTest(final String name, final Consumer public void init(final ProcessorContext context) { final T store = (T) context.getStateStore(name); - checkStateStoreMethods(store); - checker.accept(store); } @@ -235,7 +383,7 @@ public void close() { processor.init(context); } - private void checkStateStoreMethods(final StateStore store) { + private void checkGlobalStateStoreMethods(final StateStore store) { checkThrowsUnsupportedOperation(store::flush, "flush"); assertEquals(STORE_NAME, store.name()); @@ -249,6 +397,18 @@ private void checkStateStoreMethods(final StateStore store) { assertTrue(closeExecuted); } + private void checkLocalStateStoreMethods(final StateStore store) { + store.flush(); + assertTrue(flushExecuted); + + assertEquals(STORE_NAME, store.name()); + assertTrue(store.persistent()); + assertTrue(store.isOpen()); + + checkThrowsUnsupportedOperation(() -> store.init(null, null), "init"); + checkThrowsUnsupportedOperation(store::close, "init"); + } + private void checkThrowsUnsupportedOperation(final Runnable check, final String name) { try { check.run(); From 84c0e801cfeff8e62f07af6110cceaffd02d3197 Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Sun, 9 Dec 2018 13:30:34 +0300 Subject: [PATCH 3/7] KAFKA-6970: Fixing test. --- .../streams/processor/internals/ProcessorTopologyTest.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorTopologyTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorTopologyTest.java index 11050fe6f5538..14b94dadab41d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorTopologyTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorTopologyTest.java @@ -662,11 +662,6 @@ public void init(final ProcessorContext context) { public void process(final String key, final String value) { store.put(key, value); } - - @Override - public void close() { - store.close(); - } } private ProcessorSupplier define(final Processor processor) { From 54c1b9c9c90c2a1c6819779ce9e826e3b9d4af2f Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Sun, 9 Dec 2018 15:34:47 +0300 Subject: [PATCH 4/7] KAFKA-6970: Tests fixed. Also fixing wrap classes for StateStore * Use existing `AbstractStateStore implements WrappedStateStore`. * Fixing `TupleForwarder` to work with several `WrappedStateStore` that wraps each other. --- .../kstream/internals/TupleForwarder.java | 14 +- .../internals/ProcessorContextImpl.java | 158 +++++++----------- .../state/internals/WrappedStateStore.java | 2 +- 3 files changed, 72 insertions(+), 102 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java index ff3ef44894be6..fddd389507739 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java @@ -50,9 +50,17 @@ class TupleForwarder { private CachedStateStore cachedStateStore(final StateStore store) { if (store instanceof CachedStateStore) { return (CachedStateStore) store; - } else if (store instanceof WrappedStateStore - && ((WrappedStateStore) store).wrappedStore() instanceof CachedStateStore) { - return (CachedStateStore) ((WrappedStateStore) store).wrappedStore(); + } else if (store instanceof WrappedStateStore ) { + StateStore wrapped = ((WrappedStateStore) store).wrappedStore(); + + while (wrapped instanceof WrappedStateStore && !(wrapped instanceof CachedStateStore)) { + wrapped = ((WrappedStateStore) wrapped).wrappedStore(); + } + + if (!(wrapped instanceof CachedStateStore)) + return null; + + return (CachedStateStore) wrapped; } return null; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java index a862dc26fe65a..521e0c8d56918 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java @@ -38,6 +38,7 @@ import java.time.Duration; import java.util.List; +import org.apache.kafka.streams.state.internals.WrappedStateStore.AbstractStateStore; import static org.apache.kafka.streams.internals.ApiUtils.prepareMillisCheckFailMsgPrefix; @@ -206,69 +207,47 @@ public long streamTime() { return streamTimeSupplier.get(); } - private abstract static class StateStoreReadOnlyDecorator implements StateStore { + private abstract static class StateStoreReadOnlyDecorator extends AbstractStateStore { static final String ERROR_MESSAGE = "Global store is read only"; - final T underlying; - - StateStoreReadOnlyDecorator(final T underlying) { - this.underlying = underlying; - } - - @Override - public String name() { - return underlying.name(); + StateStoreReadOnlyDecorator(final StateStore inner) { + super(inner); } - @Override - public void init(final ProcessorContext context, final StateStore root) { - underlying.init(context, root); + @SuppressWarnings("unchecked") + T getInner() { + return (T) wrappedStore(); } @Override public void flush() { throw new UnsupportedOperationException(ERROR_MESSAGE); } - - @Override - public void close() { - underlying.close(); - } - - @Override - public boolean persistent() { - return underlying.persistent(); - } - - @Override - public boolean isOpen() { - return underlying.isOpen(); - } } private static class KeyValueStoreReadOnlyDecorator extends StateStoreReadOnlyDecorator> implements KeyValueStore { - KeyValueStoreReadOnlyDecorator(final KeyValueStore underlying) { - super(underlying); + KeyValueStoreReadOnlyDecorator(final KeyValueStore inner) { + super(inner); } @Override public V get(final K key) { - return underlying.get(key); + return getInner().get(key); } @Override public KeyValueIterator range(final K from, final K to) { - return underlying.range(from, to); + return getInner().range(from, to); } @Override public KeyValueIterator all() { - return underlying.all(); + return getInner().all(); } @Override public long approximateNumEntries() { - return underlying.approximateNumEntries(); + return getInner().approximateNumEntries(); } @Override @@ -293,8 +272,8 @@ public V delete(final K key) { } private static class WindowStoreReadOnlyDecorator extends StateStoreReadOnlyDecorator> implements WindowStore { - WindowStoreReadOnlyDecorator(final WindowStore underlying) { - super(underlying); + WindowStoreReadOnlyDecorator(final WindowStore inner) { + super(inner); } @Override @@ -309,46 +288,46 @@ public void put(final K key, final V value, final long windowStartTimestamp) { @Override public V fetch(final K key, final long time) { - return underlying.fetch(key, time); + return getInner().fetch(key, time); } @Deprecated @Override public WindowStoreIterator fetch(final K key, final long timeFrom, final long timeTo) { - return underlying.fetch(key, timeFrom, timeTo); + return getInner().fetch(key, timeFrom, timeTo); } @Deprecated @Override public KeyValueIterator, V> fetch(final K from, final K to, final long timeFrom, final long timeTo) { - return underlying.fetch(from, to, timeFrom, timeTo); + return getInner().fetch(from, to, timeFrom, timeTo); } @Override public KeyValueIterator, V> all() { - return underlying.all(); + return getInner().all(); } @Deprecated @Override public KeyValueIterator, V> fetchAll(final long timeFrom, final long timeTo) { - return underlying.fetchAll(timeFrom, timeTo); + return getInner().fetchAll(timeFrom, timeTo); } } private static class SessionStoreReadOnlyDecorator extends StateStoreReadOnlyDecorator> implements SessionStore { - SessionStoreReadOnlyDecorator(final SessionStore underlying) { - super(underlying); + SessionStoreReadOnlyDecorator(final SessionStore inner) { + super(inner); } @Override public KeyValueIterator, AGG> findSessions(final K key, final long earliestSessionEndTime, final long latestSessionStartTime) { - return underlying.findSessions(key, earliestSessionEndTime, latestSessionStartTime); + return getInner().findSessions(key, earliestSessionEndTime, latestSessionStartTime); } @Override public KeyValueIterator, AGG> findSessions(final K keyFrom, final K keyTo, final long earliestSessionEndTime, final long latestSessionStartTime) { - return underlying.findSessions(keyFrom, keyTo, earliestSessionEndTime, latestSessionStartTime); + return getInner().findSessions(keyFrom, keyTo, earliestSessionEndTime, latestSessionStartTime); } @Override @@ -363,22 +342,25 @@ public void put(final Windowed sessionKey, final AGG aggregate) { @Override public KeyValueIterator, AGG> fetch(final K key) { - return underlying.fetch(key); + return getInner().fetch(key); } @Override public KeyValueIterator, AGG> fetch(final K from, final K to) { - return underlying.fetch(from, to); + return getInner().fetch(from, to); } } - private abstract static class StateStoreReadWriteDecorator implements StateStore { + private abstract static class StateStoreReadWriteDecorator extends AbstractStateStore { static final String ERROR_MESSAGE = "This method may only be called by Kafka Streams"; - final T underlying; + StateStoreReadWriteDecorator(final T inner) { + super(inner); + } - StateStoreReadWriteDecorator(final T underlying) { - this.underlying = underlying; + @SuppressWarnings("unchecked") + T wrapped() { + return (T) super.wrappedStore(); } @Override @@ -390,151 +372,131 @@ public void init(final ProcessorContext context, final StateStore root) { public void close() { throw new UnsupportedOperationException(ERROR_MESSAGE); } - - @Override - public String name() { - return underlying.name(); - } - - @Override - public void flush() { - underlying.flush(); - } - - @Override - public boolean persistent() { - return underlying.persistent(); - } - - @Override - public boolean isOpen() { - return underlying.isOpen(); - } } private static class KeyValueStoreReadWriteDecorator extends StateStoreReadWriteDecorator> implements KeyValueStore { - KeyValueStoreReadWriteDecorator(final KeyValueStore underlying) { - super(underlying); + KeyValueStoreReadWriteDecorator(final KeyValueStore inner) { + super(inner); } @Override public V get(final K key) { - return underlying.get(key); + return wrapped().get(key); } @Override public KeyValueIterator range(final K from, final K to) { - return underlying.range(from, to); + return wrapped().range(from, to); } @Override public KeyValueIterator all() { - return underlying.all(); + return wrapped().all(); } @Override public long approximateNumEntries() { - return underlying.approximateNumEntries(); + return wrapped().approximateNumEntries(); } @Override public void put(final K key, final V value) { - underlying.put(key, value); + wrapped().put(key, value); } @Override public V putIfAbsent(final K key, final V value) { - return underlying.putIfAbsent(key, value); + return wrapped().putIfAbsent(key, value); } @Override public void putAll(final List> entries) { - underlying.putAll(entries); + wrapped().putAll(entries); } @Override public V delete(final K key) { - return underlying.delete(key); + return wrapped().delete(key); } } private static class WindowStoreReadWriteDecorator extends StateStoreReadWriteDecorator> implements WindowStore { - WindowStoreReadWriteDecorator(final WindowStore underlying) { - super(underlying); + WindowStoreReadWriteDecorator(final WindowStore inner) { + super(inner); } @Override public void put(final K key, final V value) { - underlying.put(key, value); + wrapped().put(key, value); } @Override public void put(final K key, final V value, final long windowStartTimestamp) { - underlying.put(key, value, windowStartTimestamp); + wrapped().put(key, value, windowStartTimestamp); } @Override public V fetch(final K key, final long time) { - return underlying.fetch(key, time); + return wrapped().fetch(key, time); } @Deprecated @Override public WindowStoreIterator fetch(final K key, final long timeFrom, final long timeTo) { - return underlying.fetch(key, timeFrom, timeTo); + return wrapped().fetch(key, timeFrom, timeTo); } @Deprecated @Override public KeyValueIterator, V> fetch(final K from, final K to, final long timeFrom, final long timeTo) { - return underlying.fetch(from, to, timeFrom, timeTo); + return wrapped().fetch(from, to, timeFrom, timeTo); } @Override public KeyValueIterator, V> all() { - return underlying.all(); + return wrapped().all(); } @Deprecated @Override public KeyValueIterator, V> fetchAll(final long timeFrom, final long timeTo) { - return underlying.fetchAll(timeFrom, timeTo); + return wrapped().fetchAll(timeFrom, timeTo); } } private static class SessionStoreReadWriteDecorator extends StateStoreReadWriteDecorator> implements SessionStore { - SessionStoreReadWriteDecorator(final SessionStore underlying) { - super(underlying); + SessionStoreReadWriteDecorator(final SessionStore inner) { + super(inner); } @Override public KeyValueIterator, AGG> findSessions(final K key, final long earliestSessionEndTime, final long latestSessionStartTime) { - return underlying.findSessions(key, earliestSessionEndTime, latestSessionStartTime); + return wrapped().findSessions(key, earliestSessionEndTime, latestSessionStartTime); } @Override public KeyValueIterator, AGG> findSessions(final K keyFrom, final K keyTo, final long earliestSessionEndTime, final long latestSessionStartTime) { - return underlying.findSessions(keyFrom, keyTo, earliestSessionEndTime, latestSessionStartTime); + return wrapped().findSessions(keyFrom, keyTo, earliestSessionEndTime, latestSessionStartTime); } @Override public void remove(final Windowed sessionKey) { - underlying.remove(sessionKey); + wrapped().remove(sessionKey); } @Override public void put(final Windowed sessionKey, final AGG aggregate) { - underlying.put(sessionKey, aggregate); + wrapped().put(sessionKey, aggregate); } @Override public KeyValueIterator, AGG> fetch(final K key) { - return underlying.fetch(key); + return wrapped().fetch(key); } @Override public KeyValueIterator, AGG> fetch(final K from, final K to) { - return underlying.fetch(from, to); + return wrapped().fetch(from, to); } } } diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappedStateStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappedStateStore.java index 38f966edabdb5..570c465f77fe7 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappedStateStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/WrappedStateStore.java @@ -41,7 +41,7 @@ public interface WrappedStateStore extends StateStore { abstract class AbstractStateStore implements WrappedStateStore { final StateStore innerState; - AbstractStateStore(final StateStore inner) { + protected AbstractStateStore(final StateStore inner) { this.innerState = inner; } From 970e6127541e5933a5135fee0e08aceb96bf63e7 Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Sun, 9 Dec 2018 17:54:20 +0300 Subject: [PATCH 5/7] KAFKA-6970: Style fix. --- .../apache/kafka/streams/kstream/internals/TupleForwarder.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java index fddd389507739..c91e6e507dab7 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java @@ -50,7 +50,7 @@ class TupleForwarder { private CachedStateStore cachedStateStore(final StateStore store) { if (store instanceof CachedStateStore) { return (CachedStateStore) store; - } else if (store instanceof WrappedStateStore ) { + } else if (store instanceof WrappedStateStore) { StateStore wrapped = ((WrappedStateStore) store).wrappedStore(); while (wrapped instanceof WrappedStateStore && !(wrapped instanceof CachedStateStore)) { From ba977d948e88c762e84000b9e9b465899056a479 Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Mon, 10 Dec 2018 10:05:34 +0300 Subject: [PATCH 6/7] KAFKA-6970: Code review fixes. --- .../internals/ProcessorContextImpl.java | 12 ++- .../internals/ProcessorContextImplTest.java | 91 +++++++------------ 2 files changed, 46 insertions(+), 57 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java index 521e0c8d56918..e7dd4dbc42a9c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java @@ -210,7 +210,7 @@ public long streamTime() { private abstract static class StateStoreReadOnlyDecorator extends AbstractStateStore { static final String ERROR_MESSAGE = "Global store is read only"; - StateStoreReadOnlyDecorator(final StateStore inner) { + StateStoreReadOnlyDecorator(final T inner) { super(inner); } @@ -223,6 +223,16 @@ T getInner() { public void flush() { throw new UnsupportedOperationException(ERROR_MESSAGE); } + + @Override + public void init(final ProcessorContext context, final StateStore root) { + throw new UnsupportedOperationException(ERROR_MESSAGE); + } + + @Override + public void close() { + throw new UnsupportedOperationException(ERROR_MESSAGE); + } } private static class KeyValueStoreReadOnlyDecorator extends StateStoreReadOnlyDecorator> implements KeyValueStore { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java index bb6456ba9e24c..f956e0ed5a40a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorContextImplTest.java @@ -58,8 +58,6 @@ public class ProcessorContextImplTest { private static final long VAL = 42L; private static final String STORE_NAME = "underlying-store"; - private boolean initExecuted; - private boolean closeExecuted; private boolean flushExecuted; private boolean putExecuted; private boolean putIfAbsentExecuted; @@ -76,8 +74,6 @@ public class ProcessorContextImplTest { @Before public void setup() { - initExecuted = false; - closeExecuted = false; flushExecuted = false; putExecuted = false; putIfAbsentExecuted = false; @@ -128,14 +124,15 @@ public void setup() { } @Test - public void testGlobalKeyValueStore() { + public void globalKeyValueStoreShouldBeReadOnly() { doTest("GlobalKeyValueStore", (Consumer>) store -> { - checkGlobalStateStoreMethods(store); + verifyStoreCannotBeInitializedOrClosed(store); - checkThrowsUnsupportedOperation(() -> store.put("1", 1L), "put"); - checkThrowsUnsupportedOperation(() -> store.putIfAbsent("1", 1L), "putIfAbsent"); - checkThrowsUnsupportedOperation(() -> store.putAll(Collections.emptyList()), "putAll"); - checkThrowsUnsupportedOperation(() -> store.delete("1"), "delete"); + checkThrowsUnsupportedOperation(store::flush, "flush()"); + checkThrowsUnsupportedOperation(() -> store.put("1", 1L), "put()"); + checkThrowsUnsupportedOperation(() -> store.putIfAbsent("1", 1L), "putIfAbsent()"); + checkThrowsUnsupportedOperation(() -> store.putAll(Collections.emptyList()), "putAll()"); + checkThrowsUnsupportedOperation(() -> store.delete("1"), "delete()"); assertEquals((Long) VAL, store.get(KEY)); assertEquals(rangeIter, store.range("one", "two")); @@ -145,12 +142,13 @@ public void testGlobalKeyValueStore() { } @Test - public void testGlobalWindowStore() { + public void globalWindowStoreShouldBeReadOnly() { doTest("GlobalWindowStore", (Consumer>) store -> { - checkGlobalStateStoreMethods(store); + verifyStoreCannotBeInitializedOrClosed(store); - checkThrowsUnsupportedOperation(() -> store.put("1", 1L, 1L), "put"); - checkThrowsUnsupportedOperation(() -> store.put("1", 1L), "put"); + checkThrowsUnsupportedOperation(store::flush, "flush()"); + checkThrowsUnsupportedOperation(() -> store.put("1", 1L, 1L), "put()"); + checkThrowsUnsupportedOperation(() -> store.put("1", 1L), "put()"); assertEquals(iters.get(0), store.fetchAll(0L, 0L)); assertEquals(windowStoreIter, store.fetch(KEY, 0L, 1L)); @@ -161,12 +159,13 @@ public void testGlobalWindowStore() { } @Test - public void testGlobalSessionStore() { + public void globalSessionStoreShouldBeReadOnly() { doTest("GlobalSessionStore", (Consumer>) store -> { - checkGlobalStateStoreMethods(store); + verifyStoreCannotBeInitializedOrClosed(store); - checkThrowsUnsupportedOperation(() -> store.remove(null), "remove"); - checkThrowsUnsupportedOperation(() -> store.put(null, null), "put"); + checkThrowsUnsupportedOperation(store::flush, "flush()"); + checkThrowsUnsupportedOperation(() -> store.remove(null), "remove()"); + checkThrowsUnsupportedOperation(() -> store.put(null, null), "put()"); assertEquals(iters.get(3), store.findSessions(KEY, 1L, 2L)); assertEquals(iters.get(4), store.findSessions(KEY, KEY, 1L, 2L)); @@ -176,9 +175,12 @@ public void testGlobalSessionStore() { } @Test - public void testLocalKeyValueStore() { + public void localKeyValueStoreShouldNotAllowInitOrClose() { doTest("LocalKeyValueStore", (Consumer>) store -> { - checkLocalStateStoreMethods(store); + verifyStoreCannotBeInitializedOrClosed(store); + + store.flush(); + assertTrue(flushExecuted); store.put("1", 1L); assertTrue(putExecuted); @@ -200,9 +202,12 @@ public void testLocalKeyValueStore() { } @Test - public void testLocalWindowStore() { + public void localWindowStoreShouldNotAllowInitOrClose() { doTest("LocalWindowStore", (Consumer>) store -> { - checkLocalStateStoreMethods(store); + verifyStoreCannotBeInitializedOrClosed(store); + + store.flush(); + assertTrue(flushExecuted); store.put("1", 1L); assertTrue(putExecuted); @@ -219,9 +224,12 @@ public void testLocalWindowStore() { } @Test - public void testLocalSessionStore() { + public void localSessionStoreShouldNotAllowInitOrClose() { doTest("LocalSessionStore", (Consumer>) store -> { - checkLocalStateStoreMethods(store); + verifyStoreCannotBeInitializedOrClosed(store); + + store.flush(); + assertTrue(flushExecuted); store.remove(null); assertTrue(removeExecuted); @@ -339,18 +347,6 @@ private void initStateStoreMock(final StateStore stateStore) { expect(stateStore.persistent()).andReturn(true); expect(stateStore.isOpen()).andReturn(true); - stateStore.init(null, null); - expectLastCall().andAnswer(() -> { - initExecuted = true; - return null; - }); - - stateStore.close(); - expectLastCall().andAnswer(() -> { - closeExecuted = true; - return null; - }); - stateStore.flush(); expectLastCall().andAnswer(() -> { flushExecuted = true; @@ -383,30 +379,13 @@ public void close() { processor.init(context); } - private void checkGlobalStateStoreMethods(final StateStore store) { - checkThrowsUnsupportedOperation(store::flush, "flush"); - - assertEquals(STORE_NAME, store.name()); - assertTrue(store.persistent()); - assertTrue(store.isOpen()); - - store.init(null, null); - assertTrue(initExecuted); - - store.close(); - assertTrue(closeExecuted); - } - - private void checkLocalStateStoreMethods(final StateStore store) { - store.flush(); - assertTrue(flushExecuted); - + private void verifyStoreCannotBeInitializedOrClosed(final StateStore store) { assertEquals(STORE_NAME, store.name()); assertTrue(store.persistent()); assertTrue(store.isOpen()); - checkThrowsUnsupportedOperation(() -> store.init(null, null), "init"); - checkThrowsUnsupportedOperation(store::close, "init"); + checkThrowsUnsupportedOperation(() -> store.init(null, null), "init()"); + checkThrowsUnsupportedOperation(store::close, "close()"); } private void checkThrowsUnsupportedOperation(final Runnable check, final String name) { From 6e2ba437831d97e39916c5e6f99380ccde4a2168 Mon Sep 17 00:00:00 2001 From: Nikolay Izhikov Date: Mon, 10 Dec 2018 19:11:39 +0300 Subject: [PATCH 7/7] KAFKA-6970: Code review fixes. --- .../apache/kafka/streams/kstream/internals/TupleForwarder.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java index c91e6e507dab7..99ba0f6ce06b7 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/TupleForwarder.java @@ -57,8 +57,9 @@ private CachedStateStore cachedStateStore(final StateStore store) { wrapped = ((WrappedStateStore) wrapped).wrappedStore(); } - if (!(wrapped instanceof CachedStateStore)) + if (!(wrapped instanceof CachedStateStore)) { return null; + } return (CachedStateStore) wrapped; }