From 295db52e69e657e554b6822d061a7daccc76a6b4 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Apr 2019 16:07:10 -0700 Subject: [PATCH 1/3] Check for null result before deserializing, delegate to underlying fetch() --- .../state/internals/MeteredSessionStore.java | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java index 4631601b12a9f..a3a810a10a8b1 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java @@ -151,22 +151,23 @@ public void remove(final Windowed sessionKey) { @Override public V fetchSession(final K key, final long startTime, final long endTime) { Objects.requireNonNull(key, "key cannot be null"); - final V value; final Bytes bytesKey = keyBytes(key); final long startNs = time.nanoseconds(); try { - value = serdes.valueFrom(wrapped().fetchSession(bytesKey, startTime, endTime)); + final byte[] result = wrapped().fetchSession(bytesKey, startTime, endTime); + if (result == null) { + return null; + } + return serdes.valueFrom(result); } finally { metrics.recordLatency(flushTime, startNs, time.nanoseconds()); } - - return value; } @Override public KeyValueIterator, V> fetch(final K key) { Objects.requireNonNull(key, "key cannot be null"); - return findSessions(key, 0, Long.MAX_VALUE); + return fetch(key); } @Override @@ -174,7 +175,7 @@ public KeyValueIterator, V> fetch(final K from, final K to) { Objects.requireNonNull(from, "from cannot be null"); Objects.requireNonNull(to, "to cannot be null"); - return findSessions(from, to, 0, Long.MAX_VALUE); + return fetch(from, to); } @Override From a852193bc1d5d9d818a87215acb6a689c1f57068 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Apr 2019 17:38:19 -0700 Subject: [PATCH 2/3] Fixed bug --- .../state/internals/MeteredSessionStore.java | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java index a3a810a10a8b1..94b004e330e75 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/MeteredSessionStore.java @@ -167,7 +167,12 @@ public V fetchSession(final K key, final long startTime, final long endTime) { @Override public KeyValueIterator, V> fetch(final K key) { Objects.requireNonNull(key, "key cannot be null"); - return fetch(key); + return new MeteredWindowedKeyValueIterator<>( + wrapped().fetch(keyBytes(key)), + fetchTime, + metrics, + serdes, + time); } @Override @@ -175,7 +180,12 @@ public KeyValueIterator, V> fetch(final K from, final K to) { Objects.requireNonNull(from, "from cannot be null"); Objects.requireNonNull(to, "to cannot be null"); - return fetch(from, to); + return new MeteredWindowedKeyValueIterator<>( + wrapped().fetch(keyBytes(from), keyBytes(to)), + fetchTime, + metrics, + serdes, + time); } @Override From 661d80498dd6a1dffc3270a2ef5164159ea097a3 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 15 Apr 2019 14:39:28 -0700 Subject: [PATCH 3/3] Fixed failing test and added new ones to test null return on get/fetchSession --- .../state/internals/MeteredKeyValueStoreTest.java | 9 +++++++++ .../state/internals/MeteredSessionStoreTest.java | 13 +++++++++++-- .../state/internals/MeteredWindowStoreTest.java | 2 +- 3 files changed, 21 insertions(+), 3 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java index b8fc88e62f82d..5cbe95cea2d37 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredKeyValueStoreTest.java @@ -55,6 +55,7 @@ import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; @RunWith(EasyMockRunner.class) @@ -243,6 +244,14 @@ public void shouldSetFlushListenerOnWrappedCachingStore() { verify(cachedKeyValueStore); } + @Test + public void shouldNotThrowNullPointerExceptionIfGetReturnsNull() { + expect(inner.get(Bytes.wrap("a".getBytes()))).andReturn(null); + + init(); + assertNull(metered.get("a")); + } + @Test public void shouldNotSetFlushListenerOnWrappedNoneCachingStore() { assertFalse(metered.setFlushListener(null, false)); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java index b349f178ac257..30c382b19c1c9 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredSessionStoreTest.java @@ -57,6 +57,7 @@ import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; @RunWith(EasyMockRunner.class) @@ -172,7 +173,7 @@ public void shouldRemoveFromStoreAndRecordRemoveMetric() { @Test public void shouldFetchForKeyAndRecordFetchMetric() { - expect(inner.findSessions(Bytes.wrap(keyBytes), 0, Long.MAX_VALUE)) + expect(inner.fetch(Bytes.wrap(keyBytes))) .andReturn(new KeyValueIteratorStub<>( Collections.singleton(KeyValue.pair(windowedKeyBytes, keyBytes)).iterator())); init(); @@ -189,7 +190,7 @@ public void shouldFetchForKeyAndRecordFetchMetric() { @Test public void shouldFetchRangeFromStoreAndRecordFetchMetric() { - expect(inner.findSessions(Bytes.wrap(keyBytes), Bytes.wrap(keyBytes), 0, Long.MAX_VALUE)) + expect(inner.fetch(Bytes.wrap(keyBytes), Bytes.wrap(keyBytes))) .andReturn(new KeyValueIteratorStub<>( Collections.singleton(KeyValue.pair(windowedKeyBytes, keyBytes)).iterator())); init(); @@ -211,6 +212,14 @@ public void shouldRecordRestoreTimeOnInit() { assertTrue((Double) metric.metricValue() > 0); } + @Test + public void shouldNotThrowNullPointerExceptionIfFetchSessionReturnsNull() { + expect(inner.fetchSession(Bytes.wrap("a".getBytes()), 0, Long.MAX_VALUE)).andReturn(null); + + init(); + assertNull(metered.fetchSession("a", 0, Long.MAX_VALUE)); + } + @Test(expected = NullPointerException.class) public void shouldThrowNullPointerOnPutIfKeyIsNull() { metered.put(null, "a"); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java index 962888ae9cd41..c0ed7f634017e 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/MeteredWindowStoreTest.java @@ -178,7 +178,7 @@ public void shouldCloseUnderlyingStore() { } @Test - public void shouldNotExceptionIfFetchReturnsNull() { + public void shouldNotThrowNullPointerExceptionIfFetchReturnsNull() { expect(innerStoreMock.fetch(Bytes.wrap("a".getBytes()), 0)).andReturn(null); replay(innerStoreMock);