From 2f4ba9eca7ebfb898f9d0b445501a52761442fae Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 8 Jul 2019 12:38:55 -0700 Subject: [PATCH 1/2] Close writebatch and reorder closing of options --- .../state/internals/AbstractRocksDBSegmentedBytesStore.java | 1 + .../org/apache/kafka/streams/state/internals/RocksDBStore.java | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java index 22f3a0249a4c6..ef18d3ca28b62 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractRocksDBSegmentedBytesStore.java @@ -233,6 +233,7 @@ void restoreAllInternal(final Collection> records) { final S segment = entry.getKey(); final WriteBatch batch = entry.getValue(); segment.write(batch); + batch.close(); } } catch (final RocksDBException e) { throw new ProcessorStateException("Error restoring batch to store " + this.name, e); diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java index 0643e64a8be25..7651c97433faa 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java @@ -403,10 +403,10 @@ public synchronized void close() { } dbAccessor.close(); + db.close(); userSpecifiedOptions.close(); wOptions.close(); fOptions.close(); - db.close(); filter.close(); cache.close(); From 1203d5c28a8b8db23fb53ce5d3c5cc857b4411f6 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Thu, 11 Jul 2019 15:40:18 -0700 Subject: [PATCH 2/2] remove reording of close --- .../org/apache/kafka/streams/state/internals/RocksDBStore.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java index 7651c97433faa..0643e64a8be25 100644 --- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java +++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java @@ -403,10 +403,10 @@ public synchronized void close() { } dbAccessor.close(); - db.close(); userSpecifiedOptions.close(); wOptions.close(); fOptions.close(); + db.close(); filter.close(); cache.close();