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..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 @@ -472,6 +472,11 @@ public void shouldLogAndMeasureExpiredRecords() { assertThat(messages, hasItem("Skipping record for expired segment.")); } + @Test + public void shouldNotThrowExceptionRemovingNonexistentKey() { + sessionStore.remove(new Windowed<>("a", new SessionWindow(0, 1))); + } + @Test(expected = NullPointerException.class) public void shouldThrowNullPointerExceptionOnFindSessionsNullKey() { sessionStore.findSessions(null, 1L, 2L);