Skip to content
5 changes: 3 additions & 2 deletions core/src/main/scala/kafka/log/Log.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
47 changes: 22 additions & 25 deletions core/src/main/scala/kafka/log/LogSegment.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Comment thread
thomaskwscott marked this conversation as resolved.
Outdated

/* 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 */
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -338,22 +338,21 @@ 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()
ensureOffsetInRange(batch.lastOffset)

// 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()
Expand All @@ -378,23 +377,21 @@ 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
}

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

Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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)
Expand Down