From 02ab0e462e4067a521a26e07d3cf8378a5aa0145 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Tue, 23 Oct 2018 20:28:45 -0400 Subject: [PATCH 1/2] KAFKA-7534: Error in flush calling close may prevent underlying store from closing --- .../state/internals/CachingKeyValueStore.java | 9 ++++++--- .../internals/CachingKeyValueStoreTest.java | 20 ++++++++++++++++++- 2 files changed, 25 insertions(+), 4 deletions(-) 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 a6a24ea098fab..a54ae1e0ed5f3 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 @@ -129,9 +129,12 @@ public void flush() { @Override public void close() { - flush(); - underlying.close(); - cache.close(cacheName); + try { + flush(); + } finally { + underlying.close(); + cache.close(cacheName); + } } @Override diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingKeyValueStoreTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingKeyValueStoreTest.java index ae6bded1cc2e3..403478ee60569 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingKeyValueStoreTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/CachingKeyValueStoreTest.java @@ -34,6 +34,7 @@ import org.apache.kafka.streams.state.StoreBuilder; import org.apache.kafka.streams.state.Stores; import org.apache.kafka.test.InternalMockProcessorContext; +import org.easymock.EasyMock; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -103,6 +104,22 @@ protected KeyValueStore createKeyValueStore(final ProcessorContext return store; } + @Test + public void shouldCloseAfterErrorWithFlush() { + try { + cache = EasyMock.niceMock(ThreadCache.class); + context = new InternalMockProcessorContext(null, null, null, (RecordCollector) null, cache); + context.setRecordContext(new ProcessorRecordContext(10, 0, 0, topic, null)); + store.init(context, null); + cache.flush("0_0-store"); + EasyMock.expectLastCall().andThrow(new NullPointerException("Simulating an error on flush")); + EasyMock.replay(cache); + store.close(); + } catch (final NullPointerException npe) { + assertFalse(underlyingStore.isOpen()); + } + } + @Test public void shouldPutGetToFromCache() { store.put(bytesKey("key"), bytesValue("value")); @@ -274,7 +291,8 @@ public void shouldThrowNullPointerExceptionOnPutAllWithNullKey() { try { store.putAll(entries); fail("Should have thrown NullPointerException while putAll null key"); - } catch (final NullPointerException e) { } + } catch (final NullPointerException e) { + } } @Test From 810a42800fc022427c14818fcdf669d9e6362245 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Tue, 23 Oct 2018 21:48:44 -0400 Subject: [PATCH 2/2] KAFKA-7534: use nested try to ensure cache closed as well. --- .../streams/state/internals/CachingKeyValueStore.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) 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 a54ae1e0ed5f3..8d9b20734c3e5 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 @@ -132,8 +132,11 @@ public void close() { try { flush(); } finally { - underlying.close(); - cache.close(cacheName); + try { + underlying.close(); + } finally { + cache.close(cacheName); + } } }