From c42cdd887af0d4806c6f7c502746ef315cf09722 Mon Sep 17 00:00:00 2001 From: Ismael Juma Date: Wed, 15 Apr 2020 09:16:46 -0700 Subject: [PATCH] MINOR: Use streaming iterator with decompression buffer when building offset map This makes it consistent with the `filterTo` methods. --- .../src/main/scala/kafka/log/LogCleaner.scala | 19 +++++++++++-------- 1 file changed, 11 insertions(+), 8 deletions(-) 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() } }