diff --git a/core/src/main/scala/kafka/log/LogCleaner.scala b/core/src/main/scala/kafka/log/LogCleaner.scala index 2a20292b8211b..2ce1a8546007c 100644 --- a/core/src/main/scala/kafka/log/LogCleaner.scala +++ b/core/src/main/scala/kafka/log/LogCleaner.scala @@ -938,15 +938,18 @@ private[log] class Cleaner(val id: Int, // Note that abort markers are supported in v2 and above, which means count is defined. stats.indexMessagesRead(batch.countOrNull) } else { - for (record <- batch.asScala) { - if (record.hasKey && record.offset >= startOffset) { - if (map.size < maxDesiredMapSize) - map.put(record.key, record.offset) - else - return true + val recordsIterator = batch.streamingIterator(decompressionBufferSupplier) + try { + for (record <- recordsIterator.asScala) { + if (record.hasKey && record.offset >= startOffset) { + if (map.size < maxDesiredMapSize) + map.put(record.key, record.offset) + else + return true + } + stats.indexMessagesRead(1) } - stats.indexMessagesRead(1) - } + } finally recordsIterator.close() } }