From b2bb961f1c1d258673685a6a3b14e00ef7bd72cb Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 29 Jul 2019 16:28:28 -0700 Subject: [PATCH 1/2] fix bug --- .../streams/state/internals/InMemorySessionStore.java | 8 ++++++++ .../streams/state/internals/SessionBytesStoreTest.java | 5 +++++ 2 files changed, 13 insertions(+) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java index f3b85657278dc..ebe9878fe2721 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/InMemorySessionStore.java @@ -120,7 +120,15 @@ public void put(final Windowed sessionKey, final byte[] aggregate) { @Override public void remove(final Windowed sessionKey) { final ConcurrentNavigableMap> keyMap = endTimeMap.get(sessionKey.window().end()); + if (keyMap == null) { + return; + } + final ConcurrentNavigableMap startTimeMap = keyMap.get(sessionKey.key()); + if (startTimeMap == null) { + return; + } + startTimeMap.remove(sessionKey.window().start()); if (startTimeMap.isEmpty()) { diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java index a6cef81b3777e..beee6c04fd6cc 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java @@ -472,6 +472,11 @@ public void shouldLogAndMeasureExpiredRecords() { assertThat(messages, hasItem("Skipping record for expired segment.")); } + @Test + public void shouldNotThrowRemovingNonexistentKey() { + sessionStore.put(new Windowed<>("a", new SessionWindow(0, 1)), 0L); + } + @Test(expected = NullPointerException.class) public void shouldThrowNullPointerExceptionOnFindSessionsNullKey() { sessionStore.findSessions(null, 1L, 2L); From 23303c20f3c74afb97767c5941d2ffffcee450cd Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Tue, 30 Jul 2019 10:48:39 -0700 Subject: [PATCH 2/2] use remove in test --- .../kafka/streams/state/internals/SessionBytesStoreTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java index beee6c04fd6cc..9b29e8b9ff607 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/SessionBytesStoreTest.java @@ -473,8 +473,8 @@ public void shouldLogAndMeasureExpiredRecords() { } @Test - public void shouldNotThrowRemovingNonexistentKey() { - sessionStore.put(new Windowed<>("a", new SessionWindow(0, 1)), 0L); + public void shouldNotThrowExceptionRemovingNonexistentKey() { + sessionStore.remove(new Windowed<>("a", new SessionWindow(0, 1))); } @Test(expected = NullPointerException.class)