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);