From 9c7a7c0224b86704337ece1bf971709c63852dcf Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Fri, 2 Jul 2021 13:05:35 +0100 Subject: [PATCH 01/10] KAFKA-12981 refactor LogSegment.offsetOfMaxTimestampSoFar and LogSegment.maxTimestampSoFar to single tuple to ensure consistent update/read --- core/src/main/scala/kafka/log/Log.scala | 5 +- .../src/main/scala/kafka/log/LogSegment.scala | 47 +++++++++---------- 2 files changed, 25 insertions(+), 27 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index 83fa50180c0a0..9c617acc0419d 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -1345,8 +1345,9 @@ class Log(@volatile private var _dir: File, val latestTimestampSegment = segmentsCopy.maxBy(_.maxTimestampSoFar) val latestEpochOpt = leaderEpochCache.flatMap(_.latestEpoch).map(_.asInstanceOf[Integer]) val epochOptional = Optional.ofNullable(latestEpochOpt.orNull) - Some(new TimestampAndOffset(latestTimestampSegment.maxTimestampSoFar, - latestTimestampSegment.offsetOfMaxTimestampSoFar, + val latestTimestampAndOffset = latestTimestampSegment.maxTimestampAndOffsetSoFar + Some(new TimestampAndOffset(latestTimestampAndOffset._1, + latestTimestampAndOffset._2, epochOptional)) } else { // Cache to avoid race conditions. `toBuffer` is faster than most alternatives and provides diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 37882ffa52592..68631bff9510d 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -99,21 +99,22 @@ class LogSegment private[log] (val log: FileRecords, // volatile for LogCleaner to see the update @volatile private var rollingBasedTimestamp: Option[Long] = None + /* The maximum timestamp and offset we see so far */ + @volatile private var _maxTimestampAndOffsetSoFar: (Option[Long], Option[Long]) = (None, None) + def maxTimestampAndOffsetSoFar_= (timestampAndOffset: (Long, Long)) : Unit = _maxTimestampAndOffsetSoFar = (Some(timestampAndOffset._1), Some(timestampAndOffset._2)) + def maxTimestampAndOffsetSoFar: (Long,Long) = { + if (_maxTimestampAndOffsetSoFar._1.isEmpty || _maxTimestampAndOffsetSoFar._2.isEmpty) + _maxTimestampAndOffsetSoFar = (Some(timeIndex.lastEntry.timestamp), Some(timeIndex.lastEntry.offset)) + (_maxTimestampAndOffsetSoFar._1.get, _maxTimestampAndOffsetSoFar._2.get) + } + /* The maximum timestamp we see so far */ - @volatile private var _maxTimestampSoFar: Option[Long] = None - def maxTimestampSoFar_=(timestamp: Long): Unit = _maxTimestampSoFar = Some(timestamp) def maxTimestampSoFar: Long = { - if (_maxTimestampSoFar.isEmpty) - _maxTimestampSoFar = Some(timeIndex.lastEntry.timestamp) - _maxTimestampSoFar.get + maxTimestampAndOffsetSoFar._1 } - @volatile private var _offsetOfMaxTimestampSoFar: Option[Long] = None - def offsetOfMaxTimestampSoFar_=(offset: Long): Unit = _offsetOfMaxTimestampSoFar = Some(offset) def offsetOfMaxTimestampSoFar: Long = { - if (_offsetOfMaxTimestampSoFar.isEmpty) - _offsetOfMaxTimestampSoFar = Some(timeIndex.lastEntry.offset) - _offsetOfMaxTimestampSoFar.get + maxTimestampAndOffsetSoFar._2 } /* Return the size in bytes of this log segment */ @@ -158,13 +159,12 @@ class LogSegment private[log] (val log: FileRecords, trace(s"Appended $appendedBytes to ${log.file} at end offset $largestOffset") // Update the in memory max timestamp and corresponding offset. if (largestTimestamp > maxTimestampSoFar) { - maxTimestampSoFar = largestTimestamp - offsetOfMaxTimestampSoFar = shallowOffsetOfMaxTimestamp + maxTimestampAndOffsetSoFar = (largestTimestamp, shallowOffsetOfMaxTimestamp) } // append an entry to the index (if needed) if (bytesSinceLastIndexEntry > indexIntervalBytes) { offsetIndex.append(largestOffset, physicalPosition) - timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2) bytesSinceLastIndexEntry = 0 } bytesSinceLastIndexEntry += records.sizeInBytes @@ -338,7 +338,7 @@ class LogSegment private[log] (val log: FileRecords, txnIndex.reset() var validBytes = 0 var lastIndexEntry = 0 - maxTimestampSoFar = RecordBatch.NO_TIMESTAMP + maxTimestampAndOffsetSoFar = (RecordBatch.NO_TIMESTAMP, 0L) try { for (batch <- log.batches.asScala) { batch.ensureValid() @@ -346,14 +346,13 @@ class LogSegment private[log] (val log: FileRecords, // The max timestamp is exposed at the batch level, so no need to iterate the records if (batch.maxTimestamp > maxTimestampSoFar) { - maxTimestampSoFar = batch.maxTimestamp - offsetOfMaxTimestampSoFar = batch.lastOffset + maxTimestampAndOffsetSoFar = (batch.maxTimestamp, batch.lastOffset) } // Build offset index if (validBytes - lastIndexEntry > indexIntervalBytes) { offsetIndex.append(batch.lastOffset, validBytes) - timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2) lastIndexEntry = validBytes } validBytes += batch.sizeInBytes() @@ -378,7 +377,7 @@ class LogSegment private[log] (val log: FileRecords, log.truncateTo(validBytes) offsetIndex.trimToValidSize() // A normally closed segment always appends the biggest timestamp ever seen into log segment, we do this as well. - timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar, skipFullCheck = true) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2, skipFullCheck = true) timeIndex.trimToValidSize() truncated } @@ -386,15 +385,13 @@ class LogSegment private[log] (val log: FileRecords, private def loadLargestTimestamp(): Unit = { // Get the last time index entry. If the time index is empty, it will return (-1, baseOffset) val lastTimeIndexEntry = timeIndex.lastEntry - maxTimestampSoFar = lastTimeIndexEntry.timestamp - offsetOfMaxTimestampSoFar = lastTimeIndexEntry.offset + maxTimestampAndOffsetSoFar = (lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) val offsetPosition = offsetIndex.lookup(lastTimeIndexEntry.offset) // Scan the rest of the messages to see if there is a larger timestamp after the last time index entry. val maxTimestampOffsetAfterLastEntry = log.largestTimestampAfter(offsetPosition.position) if (maxTimestampOffsetAfterLastEntry.timestamp > lastTimeIndexEntry.timestamp) { - maxTimestampSoFar = maxTimestampOffsetAfterLastEntry.timestamp - offsetOfMaxTimestampSoFar = maxTimestampOffsetAfterLastEntry.offset + maxTimestampAndOffsetSoFar = (maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) } } @@ -503,7 +500,7 @@ class LogSegment private[log] (val log: FileRecords, * The time index entry appended will be used to decide when to delete the segment. */ def onBecomeInactiveSegment(): Unit = { - timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar, skipFullCheck = true) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2, skipFullCheck = true) offsetIndex.trimToValidSize() timeIndex.trimToValidSize() log.trim() @@ -584,8 +581,8 @@ class LogSegment private[log] (val log: FileRecords, * Close this log segment */ def close(): Unit = { - if (_maxTimestampSoFar.nonEmpty || _offsetOfMaxTimestampSoFar.nonEmpty) - CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar, + if (_maxTimestampAndOffsetSoFar._1.nonEmpty || _maxTimestampAndOffsetSoFar._2.nonEmpty) + CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2, skipFullCheck = true), this) CoreUtils.swallow(lazyOffsetIndex.close(), this) CoreUtils.swallow(lazyTimeIndex.close(), this) From 861a3291374414a2960656c40ec1363a9f6ed1b2 Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Fri, 2 Jul 2021 17:00:58 +0100 Subject: [PATCH 02/10] KAFKA-12981 refactored tuple to case class --- core/src/main/scala/kafka/log/Log.scala | 4 +- .../src/main/scala/kafka/log/LogSegment.scala | 43 +++++++++++-------- 2 files changed, 26 insertions(+), 21 deletions(-) diff --git a/core/src/main/scala/kafka/log/Log.scala b/core/src/main/scala/kafka/log/Log.scala index 9c617acc0419d..13f7598431e91 100644 --- a/core/src/main/scala/kafka/log/Log.scala +++ b/core/src/main/scala/kafka/log/Log.scala @@ -1346,8 +1346,8 @@ class Log(@volatile private var _dir: File, val latestEpochOpt = leaderEpochCache.flatMap(_.latestEpoch).map(_.asInstanceOf[Integer]) val epochOptional = Optional.ofNullable(latestEpochOpt.orNull) val latestTimestampAndOffset = latestTimestampSegment.maxTimestampAndOffsetSoFar - Some(new TimestampAndOffset(latestTimestampAndOffset._1, - latestTimestampAndOffset._2, + Some(new TimestampAndOffset(latestTimestampAndOffset.timestamp, + latestTimestampAndOffset.offset, epochOptional)) } else { // Cache to avoid race conditions. `toBuffer` is faster than most alternatives and provides diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 68631bff9510d..cb02cd5c43720 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -100,21 +100,20 @@ class LogSegment private[log] (val log: FileRecords, @volatile private var rollingBasedTimestamp: Option[Long] = None /* The maximum timestamp and offset we see so far */ - @volatile private var _maxTimestampAndOffsetSoFar: (Option[Long], Option[Long]) = (None, None) - def maxTimestampAndOffsetSoFar_= (timestampAndOffset: (Long, Long)) : Unit = _maxTimestampAndOffsetSoFar = (Some(timestampAndOffset._1), Some(timestampAndOffset._2)) - def maxTimestampAndOffsetSoFar: (Long,Long) = { - if (_maxTimestampAndOffsetSoFar._1.isEmpty || _maxTimestampAndOffsetSoFar._2.isEmpty) - _maxTimestampAndOffsetSoFar = (Some(timeIndex.lastEntry.timestamp), Some(timeIndex.lastEntry.offset)) - (_maxTimestampAndOffsetSoFar._1.get, _maxTimestampAndOffsetSoFar._2.get) + @volatile private var _maxTimestampAndOffsetSoFar: MaxTimestampAndOffset = MaxTimestampAndOffset.empty + def maxTimestampAndOffsetSoFar: MaxTimestampAndOffset = { + if (_maxTimestampAndOffsetSoFar == MaxTimestampAndOffset.empty) + _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(timeIndex.lastEntry.timestamp,timeIndex.lastEntry.offset) + _maxTimestampAndOffsetSoFar } /* The maximum timestamp we see so far */ def maxTimestampSoFar: Long = { - maxTimestampAndOffsetSoFar._1 + maxTimestampAndOffsetSoFar.timestamp } def offsetOfMaxTimestampSoFar: Long = { - maxTimestampAndOffsetSoFar._2 + maxTimestampAndOffsetSoFar.offset } /* Return the size in bytes of this log segment */ @@ -159,12 +158,12 @@ class LogSegment private[log] (val log: FileRecords, trace(s"Appended $appendedBytes to ${log.file} at end offset $largestOffset") // Update the in memory max timestamp and corresponding offset. if (largestTimestamp > maxTimestampSoFar) { - maxTimestampAndOffsetSoFar = (largestTimestamp, shallowOffsetOfMaxTimestamp) + _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(largestTimestamp, shallowOffsetOfMaxTimestamp) } // append an entry to the index (if needed) if (bytesSinceLastIndexEntry > indexIntervalBytes) { offsetIndex.append(largestOffset, physicalPosition) - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset) bytesSinceLastIndexEntry = 0 } bytesSinceLastIndexEntry += records.sizeInBytes @@ -338,7 +337,7 @@ class LogSegment private[log] (val log: FileRecords, txnIndex.reset() var validBytes = 0 var lastIndexEntry = 0 - maxTimestampAndOffsetSoFar = (RecordBatch.NO_TIMESTAMP, 0L) + _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset.empty try { for (batch <- log.batches.asScala) { batch.ensureValid() @@ -346,13 +345,13 @@ class LogSegment private[log] (val log: FileRecords, // The max timestamp is exposed at the batch level, so no need to iterate the records if (batch.maxTimestamp > maxTimestampSoFar) { - maxTimestampAndOffsetSoFar = (batch.maxTimestamp, batch.lastOffset) + _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(batch.maxTimestamp, batch.lastOffset) } // Build offset index if (validBytes - lastIndexEntry > indexIntervalBytes) { offsetIndex.append(batch.lastOffset, validBytes) - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset) lastIndexEntry = validBytes } validBytes += batch.sizeInBytes() @@ -377,7 +376,7 @@ class LogSegment private[log] (val log: FileRecords, log.truncateTo(validBytes) offsetIndex.trimToValidSize() // A normally closed segment always appends the biggest timestamp ever seen into log segment, we do this as well. - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2, skipFullCheck = true) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, skipFullCheck = true) timeIndex.trimToValidSize() truncated } @@ -385,13 +384,13 @@ class LogSegment private[log] (val log: FileRecords, private def loadLargestTimestamp(): Unit = { // Get the last time index entry. If the time index is empty, it will return (-1, baseOffset) val lastTimeIndexEntry = timeIndex.lastEntry - maxTimestampAndOffsetSoFar = (lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) + _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) val offsetPosition = offsetIndex.lookup(lastTimeIndexEntry.offset) // Scan the rest of the messages to see if there is a larger timestamp after the last time index entry. val maxTimestampOffsetAfterLastEntry = log.largestTimestampAfter(offsetPosition.position) if (maxTimestampOffsetAfterLastEntry.timestamp > lastTimeIndexEntry.timestamp) { - maxTimestampAndOffsetSoFar = (maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) + _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) } } @@ -500,7 +499,7 @@ class LogSegment private[log] (val log: FileRecords, * The time index entry appended will be used to decide when to delete the segment. */ def onBecomeInactiveSegment(): Unit = { - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2, skipFullCheck = true) + timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, skipFullCheck = true) offsetIndex.trimToValidSize() timeIndex.trimToValidSize() log.trim() @@ -581,8 +580,8 @@ class LogSegment private[log] (val log: FileRecords, * Close this log segment */ def close(): Unit = { - if (_maxTimestampAndOffsetSoFar._1.nonEmpty || _maxTimestampAndOffsetSoFar._2.nonEmpty) - CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampAndOffsetSoFar._1, maxTimestampAndOffsetSoFar._2, + if (_maxTimestampAndOffsetSoFar != MaxTimestampAndOffset.empty) + CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, skipFullCheck = true), this) CoreUtils.swallow(lazyOffsetIndex.close(), this) CoreUtils.swallow(lazyTimeIndex.close(), this) @@ -678,3 +677,9 @@ object LogSegment { object LogFlushStats extends KafkaMetricsGroup { val logFlushTimer = new KafkaTimer(newTimer("LogFlushRateAndTimeMs", TimeUnit.MILLISECONDS, TimeUnit.SECONDS)) } + +case class MaxTimestampAndOffset(timestamp: Long, offset: Long) + +object MaxTimestampAndOffset { + val empty = MaxTimestampAndOffset(RecordBatch.NO_TIMESTAMP, -1L) +} \ No newline at end of file From c805b8423cdd502e53b166d6be938457a2ba4b9b Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 10:40:25 +0100 Subject: [PATCH 03/10] KAFKA-12981 refactor to use existing TimestampOffset over new case class --- .../src/main/scala/kafka/log/LogSegment.scala | 26 +++++++------------ 1 file changed, 10 insertions(+), 16 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index cb02cd5c43720..b53bd2c56c4e4 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -100,10 +100,10 @@ class LogSegment private[log] (val log: FileRecords, @volatile private var rollingBasedTimestamp: Option[Long] = None /* The maximum timestamp and offset we see so far */ - @volatile private var _maxTimestampAndOffsetSoFar: MaxTimestampAndOffset = MaxTimestampAndOffset.empty - def maxTimestampAndOffsetSoFar: MaxTimestampAndOffset = { - if (_maxTimestampAndOffsetSoFar == MaxTimestampAndOffset.empty) - _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(timeIndex.lastEntry.timestamp,timeIndex.lastEntry.offset) + @volatile private var _maxTimestampAndOffsetSoFar: TimestampOffset = TimestampOffset.Unknown + def maxTimestampAndOffsetSoFar: TimestampOffset = { + if (_maxTimestampAndOffsetSoFar == TimestampOffset.Unknown) + _maxTimestampAndOffsetSoFar = TimestampOffset(timeIndex.lastEntry.timestamp,timeIndex.lastEntry.offset) _maxTimestampAndOffsetSoFar } @@ -158,7 +158,7 @@ class LogSegment private[log] (val log: FileRecords, trace(s"Appended $appendedBytes to ${log.file} at end offset $largestOffset") // Update the in memory max timestamp and corresponding offset. if (largestTimestamp > maxTimestampSoFar) { - _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(largestTimestamp, shallowOffsetOfMaxTimestamp) + _maxTimestampAndOffsetSoFar = TimestampOffset(largestTimestamp, shallowOffsetOfMaxTimestamp) } // append an entry to the index (if needed) if (bytesSinceLastIndexEntry > indexIntervalBytes) { @@ -337,7 +337,7 @@ class LogSegment private[log] (val log: FileRecords, txnIndex.reset() var validBytes = 0 var lastIndexEntry = 0 - _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset.empty + _maxTimestampAndOffsetSoFar = TimestampOffset.Unknown try { for (batch <- log.batches.asScala) { batch.ensureValid() @@ -345,7 +345,7 @@ class LogSegment private[log] (val log: FileRecords, // The max timestamp is exposed at the batch level, so no need to iterate the records if (batch.maxTimestamp > maxTimestampSoFar) { - _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(batch.maxTimestamp, batch.lastOffset) + _maxTimestampAndOffsetSoFar = TimestampOffset(batch.maxTimestamp, batch.lastOffset) } // Build offset index @@ -384,13 +384,13 @@ class LogSegment private[log] (val log: FileRecords, private def loadLargestTimestamp(): Unit = { // Get the last time index entry. If the time index is empty, it will return (-1, baseOffset) val lastTimeIndexEntry = timeIndex.lastEntry - _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) + _maxTimestampAndOffsetSoFar = TimestampOffset(lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) val offsetPosition = offsetIndex.lookup(lastTimeIndexEntry.offset) // Scan the rest of the messages to see if there is a larger timestamp after the last time index entry. val maxTimestampOffsetAfterLastEntry = log.largestTimestampAfter(offsetPosition.position) if (maxTimestampOffsetAfterLastEntry.timestamp > lastTimeIndexEntry.timestamp) { - _maxTimestampAndOffsetSoFar = MaxTimestampAndOffset(maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) + _maxTimestampAndOffsetSoFar = TimestampOffset(maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) } } @@ -580,7 +580,7 @@ class LogSegment private[log] (val log: FileRecords, * Close this log segment */ def close(): Unit = { - if (_maxTimestampAndOffsetSoFar != MaxTimestampAndOffset.empty) + if (_maxTimestampAndOffsetSoFar != TimestampOffset.Unknown) CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, skipFullCheck = true), this) CoreUtils.swallow(lazyOffsetIndex.close(), this) @@ -676,10 +676,4 @@ object LogSegment { object LogFlushStats extends KafkaMetricsGroup { val logFlushTimer = new KafkaTimer(newTimer("LogFlushRateAndTimeMs", TimeUnit.MILLISECONDS, TimeUnit.SECONDS)) -} - -case class MaxTimestampAndOffset(timestamp: Long, offset: Long) - -object MaxTimestampAndOffset { - val empty = MaxTimestampAndOffset(RecordBatch.NO_TIMESTAMP, -1L) } \ No newline at end of file From 86148271567c58e61a2a7d29731be06d5f717c11 Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 14:35:11 +0100 Subject: [PATCH 04/10] KAFKA-12981 fixes per pr review --- core/src/main/scala/kafka/log/LogSegment.scala | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index b53bd2c56c4e4..4e952df9df626 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -101,9 +101,10 @@ class LogSegment private[log] (val log: FileRecords, /* The maximum timestamp and offset we see so far */ @volatile private var _maxTimestampAndOffsetSoFar: TimestampOffset = TimestampOffset.Unknown + def maxTimestampAndOffsetSoFar_= (timestampOffset: TimestampOffset) : Unit = _maxTimestampAndOffsetSoFar = timestampOffset def maxTimestampAndOffsetSoFar: TimestampOffset = { if (_maxTimestampAndOffsetSoFar == TimestampOffset.Unknown) - _maxTimestampAndOffsetSoFar = TimestampOffset(timeIndex.lastEntry.timestamp,timeIndex.lastEntry.offset) + _maxTimestampAndOffsetSoFar = timeIndex.lastEntry _maxTimestampAndOffsetSoFar } @@ -158,12 +159,12 @@ class LogSegment private[log] (val log: FileRecords, trace(s"Appended $appendedBytes to ${log.file} at end offset $largestOffset") // Update the in memory max timestamp and corresponding offset. if (largestTimestamp > maxTimestampSoFar) { - _maxTimestampAndOffsetSoFar = TimestampOffset(largestTimestamp, shallowOffsetOfMaxTimestamp) + maxTimestampAndOffsetSoFar = TimestampOffset(largestTimestamp, shallowOffsetOfMaxTimestamp) } // append an entry to the index (if needed) if (bytesSinceLastIndexEntry > indexIntervalBytes) { offsetIndex.append(largestOffset, physicalPosition) - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset) + timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar) bytesSinceLastIndexEntry = 0 } bytesSinceLastIndexEntry += records.sizeInBytes @@ -376,7 +377,7 @@ class LogSegment private[log] (val log: FileRecords, log.truncateTo(validBytes) offsetIndex.trimToValidSize() // A normally closed segment always appends the biggest timestamp ever seen into log segment, we do this as well. - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, skipFullCheck = true) + timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar, skipFullCheck = true) timeIndex.trimToValidSize() truncated } @@ -499,7 +500,7 @@ class LogSegment private[log] (val log: FileRecords, * The time index entry appended will be used to decide when to delete the segment. */ def onBecomeInactiveSegment(): Unit = { - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, skipFullCheck = true) + timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar, skipFullCheck = true) offsetIndex.trimToValidSize() timeIndex.trimToValidSize() log.trim() @@ -581,7 +582,7 @@ class LogSegment private[log] (val log: FileRecords, */ def close(): Unit = { if (_maxTimestampAndOffsetSoFar != TimestampOffset.Unknown) - CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset, + CoreUtils.swallow(timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar, skipFullCheck = true), this) CoreUtils.swallow(lazyOffsetIndex.close(), this) CoreUtils.swallow(lazyTimeIndex.close(), this) From 4174b344ace6a5d45fb7d91462e941642626687a Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 14:36:04 +0100 Subject: [PATCH 05/10] KAFKA-12981 fixes per pr review --- core/src/main/scala/kafka/log/LogSegment.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 4e952df9df626..74aaf712ae208 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -677,4 +677,4 @@ object LogSegment { object LogFlushStats extends KafkaMetricsGroup { val logFlushTimer = new KafkaTimer(newTimer("LogFlushRateAndTimeMs", TimeUnit.MILLISECONDS, TimeUnit.SECONDS)) -} \ No newline at end of file +} From 97df12f5f2a82769a164b27c1f1160dac2634660 Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 14:37:45 +0100 Subject: [PATCH 06/10] KAFKA-12981 fixes per pr review --- core/src/main/scala/kafka/log/LogSegment.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 74aaf712ae208..0e025e740ffec 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -352,7 +352,7 @@ class LogSegment private[log] (val log: FileRecords, // Build offset index if (validBytes - lastIndexEntry > indexIntervalBytes) { offsetIndex.append(batch.lastOffset, validBytes) - timeIndex.maybeAppend(maxTimestampAndOffsetSoFar.timestamp, maxTimestampAndOffsetSoFar.offset) + timeIndex.maybeAppend(maxTimestampSoFar, offsetOfMaxTimestampSoFar) lastIndexEntry = validBytes } validBytes += batch.sizeInBytes() From b26d763e7b0aa11b0c8cd3eef00369bf4c9be022 Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 16:38:38 +0100 Subject: [PATCH 07/10] KAFKA-12981 fixes per pr review --- core/src/main/scala/kafka/log/LogSegment.scala | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 0e025e740ffec..7ab48e7ebae06 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -338,7 +338,7 @@ class LogSegment private[log] (val log: FileRecords, txnIndex.reset() var validBytes = 0 var lastIndexEntry = 0 - _maxTimestampAndOffsetSoFar = TimestampOffset.Unknown + maxTimestampAndOffsetSoFar = TimestampOffset.Unknown try { for (batch <- log.batches.asScala) { batch.ensureValid() @@ -346,7 +346,7 @@ class LogSegment private[log] (val log: FileRecords, // The max timestamp is exposed at the batch level, so no need to iterate the records if (batch.maxTimestamp > maxTimestampSoFar) { - _maxTimestampAndOffsetSoFar = TimestampOffset(batch.maxTimestamp, batch.lastOffset) + maxTimestampAndOffsetSoFar = TimestampOffset(batch.maxTimestamp, batch.lastOffset) } // Build offset index @@ -385,13 +385,13 @@ class LogSegment private[log] (val log: FileRecords, private def loadLargestTimestamp(): Unit = { // Get the last time index entry. If the time index is empty, it will return (-1, baseOffset) val lastTimeIndexEntry = timeIndex.lastEntry - _maxTimestampAndOffsetSoFar = TimestampOffset(lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) + maxTimestampAndOffsetSoFar = TimestampOffset(lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) val offsetPosition = offsetIndex.lookup(lastTimeIndexEntry.offset) // Scan the rest of the messages to see if there is a larger timestamp after the last time index entry. val maxTimestampOffsetAfterLastEntry = log.largestTimestampAfter(offsetPosition.position) if (maxTimestampOffsetAfterLastEntry.timestamp > lastTimeIndexEntry.timestamp) { - _maxTimestampAndOffsetSoFar = TimestampOffset(maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) + maxTimestampAndOffsetSoFar = TimestampOffset(maxTimestampOffsetAfterLastEntry.timestamp, maxTimestampOffsetAfterLastEntry.offset) } } From 1ad2126913c1cda21c0763ea4fff2d5e22ac8e7f Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 17:04:57 +0100 Subject: [PATCH 08/10] Trigger Build From 33779c6083ea302e44fc4dcd8689a23845e3929f Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Mon, 5 Jul 2021 17:23:45 +0100 Subject: [PATCH 09/10] KAFKA-12981 fixes per pr review --- core/src/main/scala/kafka/log/LogSegment.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index 7ab48e7ebae06..a718267613bcf 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -385,7 +385,7 @@ class LogSegment private[log] (val log: FileRecords, private def loadLargestTimestamp(): Unit = { // Get the last time index entry. If the time index is empty, it will return (-1, baseOffset) val lastTimeIndexEntry = timeIndex.lastEntry - maxTimestampAndOffsetSoFar = TimestampOffset(lastTimeIndexEntry.timestamp, lastTimeIndexEntry.offset) + maxTimestampAndOffsetSoFar = lastTimeIndexEntry val offsetPosition = offsetIndex.lookup(lastTimeIndexEntry.offset) // Scan the rest of the messages to see if there is a larger timestamp after the last time index entry. From 9bb34bf3beb057dd7c29e26909fef67ec9a1ecb3 Mon Sep 17 00:00:00 2001 From: thomaskwscott Date: Tue, 6 Jul 2021 16:16:29 +0100 Subject: [PATCH 10/10] KAFKA-12981 fixes per pr review --- core/src/main/scala/kafka/log/LogSegment.scala | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/log/LogSegment.scala b/core/src/main/scala/kafka/log/LogSegment.scala index a718267613bcf..7ade6177642d7 100755 --- a/core/src/main/scala/kafka/log/LogSegment.scala +++ b/core/src/main/scala/kafka/log/LogSegment.scala @@ -101,7 +101,7 @@ class LogSegment private[log] (val log: FileRecords, /* The maximum timestamp and offset we see so far */ @volatile private var _maxTimestampAndOffsetSoFar: TimestampOffset = TimestampOffset.Unknown - def maxTimestampAndOffsetSoFar_= (timestampOffset: TimestampOffset) : Unit = _maxTimestampAndOffsetSoFar = timestampOffset + def maxTimestampAndOffsetSoFar_= (timestampOffset: TimestampOffset): Unit = _maxTimestampAndOffsetSoFar = timestampOffset def maxTimestampAndOffsetSoFar: TimestampOffset = { if (_maxTimestampAndOffsetSoFar == TimestampOffset.Unknown) _maxTimestampAndOffsetSoFar = timeIndex.lastEntry @@ -527,7 +527,7 @@ class LogSegment private[log] (val log: FileRecords, * segment is rolled if the difference between the current wall clock time and the segment create time exceeds the * segment rolling time. */ - def timeWaitedForRoll(now: Long, messageTimestamp: Long) : Long = { + def timeWaitedForRoll(now: Long, messageTimestamp: Long): Long = { // Load the timestamp of the first message into memory loadFirstBatchTimestamp() rollingBasedTimestamp match { @@ -539,7 +539,7 @@ class LogSegment private[log] (val log: FileRecords, /** * @return the first batch timestamp if the timestamp is available. Otherwise return Long.MaxValue */ - def getFirstBatchTimestamp() : Long = { + def getFirstBatchTimestamp(): Long = { loadFirstBatchTimestamp() rollingBasedTimestamp match { case Some(t) if t >= 0 => t