Skip to content
Merged
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
133 changes: 89 additions & 44 deletions core/src/main/scala/kafka/log/Log.scala
Original file line number Diff line number Diff line change
Expand Up @@ -795,7 +795,9 @@ class Log(@volatile private var _dir: File,
if (truncatedBytes > 0) {
// we had an invalid message, delete all remaining log
warn(s"Corruption found in segment ${segment.baseOffset}, truncating to offset ${segment.readNextOffset}")
removeAndDeleteSegments(unflushed.toList, asyncDelete = true)
removeAndDeleteSegments(unflushed.toList,
asyncDelete = true,
reason = LogRecoveryDeletion)
truncated = true
}
}
Expand All @@ -806,7 +808,9 @@ class Log(@volatile private var _dir: File,
if (logEndOffset < logStartOffset) {
warn(s"Deleting all segments because logEndOffset ($logEndOffset) is smaller than logStartOffset ($logStartOffset). " +
"This could happen if segment files were deleted from the file system.")
removeAndDeleteSegments(logSegments, asyncDelete = true)
removeAndDeleteSegments(logSegments,
asyncDelete = true,
reason = LogRecoveryDeletion)
}
}

Expand Down Expand Up @@ -1697,16 +1701,18 @@ class Log(@volatile private var _dir: File,
* (if there is one) and returns true iff it is deletable
* @return The number of segments deleted
*/
private def deleteOldSegments(predicate: (LogSegment, Option[LogSegment]) => Boolean) = {
private def deleteOldSegments(predicate: (LogSegment, Option[LogSegment]) => Boolean,
reason: SegmentDeletionReason): Int = {
lock synchronized {
val deletable = deletableSegments(predicate)
if (deletable.nonEmpty) {
deleteSegments(deletable)
} else 0
if (deletable.nonEmpty)
deleteSegments(deletable, reason)
else
0
}
}

private def deleteSegments(deletable: Iterable[LogSegment]): Int = {
private def deleteSegments(deletable: Iterable[LogSegment], reason: SegmentDeletionReason): Int = {
maybeHandleIOException(s"Error while deleting segments for $topicPartition in dir ${dir.getParent}") {
val numToDelete = deletable.size
if (numToDelete > 0) {
Expand All @@ -1716,7 +1722,7 @@ class Log(@volatile private var _dir: File,
lock synchronized {
checkIfMemoryMappedBufferClosed()
// remove the segments for lookups
removeAndDeleteSegments(deletable, asyncDelete = true)
removeAndDeleteSegments(deletable, asyncDelete = true, reason)
maybeIncrementLogStartOffset(segments.firstEntry.getValue.baseOffset, SegmentDeletion)
}
}
Expand Down Expand Up @@ -1779,57 +1785,34 @@ class Log(@volatile private var _dir: File,
if (config.retentionMs < 0) return 0
val startMs = time.milliseconds

def shouldDelete(segment: LogSegment, nextSegmentOpt: Option[LogSegment]) = {
if (startMs - segment.largestTimestamp > config.retentionMs) {
segment.largestRecordTimestamp match {
case Some(ts) =>
info(s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" retention time ${config.retentionMs}ms breach based on the largest record timestamp from the" +
s" segment, which is $ts")
case None =>
info(s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" retention time ${config.retentionMs}ms breach based on the last modified timestamp from the" +
s" segment, which is ${segment.lastModified}")
}
true
} else {
false
}
def shouldDelete(segment: LogSegment, nextSegmentOpt: Option[LogSegment]): Boolean = {
startMs - segment.largestTimestamp > config.retentionMs
}

deleteOldSegments(shouldDelete)
deleteOldSegments(shouldDelete, RetentionMsBreachDeletion)
}

private def deleteRetentionSizeBreachedSegments(): Int = {
if (config.retentionSize < 0 || size < config.retentionSize) return 0
var diff = size - config.retentionSize
def shouldDelete(segment: LogSegment, nextSegmentOpt: Option[LogSegment]) = {
def shouldDelete(segment: LogSegment, nextSegmentOpt: Option[LogSegment]): Boolean = {
if (diff - segment.size >= 0) {
diff -= segment.size
info(s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" retention size ${config.retentionSize} bytes breach. Segment size is" +
s" ${segment.size} and total log size after deletion will be ${size - diff}")
true
} else {
false
}
}

deleteOldSegments(shouldDelete)
deleteOldSegments(shouldDelete, RetentionSizeBreachDeletion)
}

private def deleteLogStartOffsetBreachedSegments(): Int = {
def shouldDelete(segment: LogSegment, nextSegmentOpt: Option[LogSegment]) = {
if (nextSegmentOpt.exists(_.baseOffset <= logStartOffset)) {
info(s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" startOffset breach. logStartOffset is $logStartOffset")
true
} else {
false
}
def shouldDelete(segment: LogSegment, nextSegmentOpt: Option[LogSegment]): Boolean = {
nextSegmentOpt.exists(_.baseOffset <= logStartOffset)
}

deleteOldSegments(shouldDelete)
deleteOldSegments(shouldDelete, StartOffsetBreachDeletion)
}

def isFuture: Boolean = dir.getName.endsWith(Log.FutureDirSuffix)
Expand Down Expand Up @@ -1920,7 +1903,7 @@ class Log(@volatile private var _dir: File,
s"=max(provided offset = $expectedNextOffset, LEO = $logEndOffset) while it already " +
s"exists and is active with size 0. Size of time index: ${activeSegment.timeIndex.entries}," +
s" size of offset index: ${activeSegment.offsetIndex.entries}.")
removeAndDeleteSegments(Seq(activeSegment), asyncDelete = true)
removeAndDeleteSegments(Seq(activeSegment), asyncDelete = true, LogRollDeletion)
} else {
throw new KafkaException(s"Trying to roll a new log segment for topic partition $topicPartition with start offset $newOffset" +
s" =max(provided offset = $expectedNextOffset, LEO = $logEndOffset) while it already exists. Existing " +
Expand Down Expand Up @@ -2052,7 +2035,7 @@ class Log(@volatile private var _dir: File,
lock synchronized {
checkIfMemoryMappedBufferClosed()
producerExpireCheck.cancel(true)
removeAndDeleteSegments(logSegments, asyncDelete = false)
removeAndDeleteSegments(logSegments, asyncDelete = false, LogDeletion)
leaderEpochCache.foreach(_.clear())
Utils.delete(dir)
// File handlers will be closed if this log is deleted
Expand Down Expand Up @@ -2103,7 +2086,7 @@ class Log(@volatile private var _dir: File,
truncateFullyAndStartAt(targetOffset)
} else {
val deletable = logSegments.filter(segment => segment.baseOffset > targetOffset)
removeAndDeleteSegments(deletable, asyncDelete = true)
removeAndDeleteSegments(deletable, asyncDelete = true, LogTruncateDeletion)
activeSegment.truncateTo(targetOffset)
updateLogEndOffset(targetOffset)
updateLogStartOffset(math.min(targetOffset, this.logStartOffset))
Expand All @@ -2126,7 +2109,7 @@ class Log(@volatile private var _dir: File,
debug(s"Truncate and start at offset $newOffset")
lock synchronized {
checkIfMemoryMappedBufferClosed()
removeAndDeleteSegments(logSegments, asyncDelete = true)
removeAndDeleteSegments(logSegments, asyncDelete = true, LogTruncateDeletion)
addSegment(LogSegment.open(dir,
baseOffset = newOffset,
config = config,
Expand Down Expand Up @@ -2227,14 +2210,17 @@ class Log(@volatile private var _dir: File,
* @param segments The log segments to schedule for deletion
* @param asyncDelete Whether the segment files should be deleted asynchronously
*/
private def removeAndDeleteSegments(segments: Iterable[LogSegment], asyncDelete: Boolean): Unit = {
private def removeAndDeleteSegments(segments: Iterable[LogSegment],
asyncDelete: Boolean,
reason: SegmentDeletionReason): Unit = {
if (segments.nonEmpty) {
lock synchronized {
// As most callers hold an iterator into the `segments` collection and `removeAndDeleteSegment` mutates it by
// removing the deleted segment, we should force materialization of the iterator here, so that results of the
// iteration remain valid and deterministic.
val toDelete = segments.toList
toDelete.foreach { segment =>
info(s"${reason.reasonString(this, segment)}")

@kowshik kowshik Aug 1, 2020

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If we passed in the deletion reason further into the deleteSegmentFiles method, it seems we can print the reason string just once for a batch of segments being deleted. And within the reason string, we can provide the reason for deleting the batch:

https://github.com/confluentinc/ce-kafka/blob/master/core/src/main/scala/kafka/log/Log.scala#L2519
https://github.com/confluentinc/ce-kafka/blob/master/core/src/main/scala/kafka/log/Log.scala#L2526

ex: info("Deleting segments due to $reason: ${segments.mkString(",")}"

where $reason provides due to retention time 1200000ms breach.

The drawback is that sometimes we can not print segment-specific information since the error message is at a batch level. But generally it may be that segment-level deletion information could bloat our server logging, so it may be better to batch the logging instead.

What are your thoughts?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

While verbose, I think having the granularity of each segment is useful. This allows us to easily reason about why a particular segment was deleted. Note that we switched from a single log per batch to a log per segment in #8850.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we think about cases where this could be an issue? Say delete records is used, causing a large number of segments to be deleted Could that trigger excessive logging?

@kowshik kowshik Aug 2, 2020

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@dhruvilshah3 Sounds good!

@dhruvilshah3 dhruvilshah3 Aug 2, 2020

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ijuma We log one message per deleted segment. This could cause temporary increase in log volume when DeleteRecords is used or when retention is lowered, for example.

Overall, we have a few options with different tradeoffs:

  1. Log a common reason per batch being deleted, including base offsets of segments being deleted. eg.
Deleting segments due to retention time 999ms breach. BaseOffsets: (0,5,...).
  1. Log a common reason per batch being deleted, including base offsets and metadata of segments. eg.
Deleting segments due to retention time 999ms breach: LogSegment(baseOffset=0, size=360, lastModifiedTime=1596387738000, largestRecordTimestamp=Some(1596387737414)),LogSegment(baseOffset=5, size=360, lastModifiedTime=1596387738000, largestRecordTimestamp=Some(1596387737414)),...
  1. Log one message per segment being deleted. This is the current behavior. eg.
Segment with base offset 0 will be deleted due to retention time 999ms breach based on the largest record timestamp from the segment, which is ...
Segment with base offset 5 will be deleted due to retention time 999ms breach based on the largest record timestamp from the segment, which is ...
...

Doing (2) may be a reasonable tradeoff. It eliminates some of the redundancy at the cost of making it to glean per segment metadata. Let me know what you think.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think it is important to capture the segment level details. In the past, we have had trouble explaining precisely why a specific segment got deleted. For example, was it because of the last modified time or the record timestamp? When users are looking to understand why data is deleted, we should be able to provide a clear answer.

My personal preference is probably 3) because I hate dealing with lists of things in log messages. Simple grepping no longer work to extract the details. Big messages also messes up console scrolling and can choke downstream systems. For segments, I am not so worried about log noise because the rate of segment creation is not that high.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this is reasonable. Logging a segment per line will make it easier for us to diagnose issues. I made the change to log a segment per line for retention-related deletions. We still batch all segments in a single line for all other deletion events, eg. log deletion, truncation, etc.

this.segments.remove(segment.baseOffset)
}
deleteSegmentFiles(toDelete, asyncDelete)
Expand Down Expand Up @@ -2686,3 +2672,62 @@ object LogMetricNames {
List(NumLogSegments, LogStartOffset, LogEndOffset, Size)
}
}

sealed trait SegmentDeletionReason {
def reasonString(log: Log, segment: LogSegment): String
}

case object RetentionMsBreachDeletion extends SegmentDeletionReason {
Comment thread
dhruvilshah3 marked this conversation as resolved.
Outdated
override def reasonString(log: Log, segment: LogSegment): String = {
val retentionMs = log.config.retentionMs
segment.largestRecordTimestamp match {
case Some(ts) =>
s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" retention time ${retentionMs}ms breach based on the largest record timestamp from the" +
s" segment, which is $ts"
case None =>
s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" retention time ${retentionMs}ms breach based on the last modified timestamp from the" +
s" segment, which is ${segment.lastModified}"
}
}
}

case object RetentionSizeBreachDeletion extends SegmentDeletionReason {
override def reasonString(log: Log, segment: LogSegment): String = {
s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" retention size ${log.config.retentionSize} bytes breach. Segment size is" +
s" ${segment.size} and total log size is ${log.size}"
}
}

case object StartOffsetBreachDeletion extends SegmentDeletionReason {
override def reasonString(log: Log, segment: LogSegment): String = {
s"Segment with base offset ${segment.baseOffset} will be deleted due to" +
s" startOffset breach. logStartOffset is ${log.logStartOffset}"
}
}

case object LogRecoveryDeletion extends SegmentDeletionReason {
override def reasonString(log: Log, segment: LogSegment): String = {
s"Segment with base offset ${segment.baseOffset} will be deleted as part of log recovery"
}
}

case object LogDeletion extends SegmentDeletionReason {
override def reasonString(log: Log, segment: LogSegment): String = {
s"Segment with base offset ${segment.baseOffset} will be deleted as the corresponding log has been deleted"
}
}

case object LogTruncateDeletion extends SegmentDeletionReason {
override def reasonString(log: Log, segment: LogSegment): String = {
s"Segment with base offset ${segment.baseOffset} will be deleted as part of log truncation"
}
}

case object LogRollDeletion extends SegmentDeletionReason {
override def reasonString(log: Log, segment: LogSegment): String = {
s"Segment with base offset ${segment.baseOffset} will be deleted as part of log roll"
}
}