-
Notifications
You must be signed in to change notification settings - Fork 15.4k
MINOR: Add log identifier/prefix printing in Log layer static functions #10742
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 3 commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -557,7 +557,7 @@ class Log(@volatile private var _dir: File, | |
| } | ||
|
|
||
| private def initializeLeaderEpochCache(): Unit = lock synchronized { | ||
| leaderEpochCache = Log.maybeCreateLeaderEpochCache(dir, topicPartition, logDirFailureChannel, recordVersion) | ||
| leaderEpochCache = Log.maybeCreateLeaderEpochCache(dir, topicPartition, logDirFailureChannel, recordVersion, logIdent) | ||
| } | ||
|
|
||
| private def updateLogEndOffset(offset: Long): Unit = { | ||
|
|
@@ -592,7 +592,7 @@ class Log(@volatile private var _dir: File, | |
| producerStateManager: ProducerStateManager): Unit = lock synchronized { | ||
| checkIfMemoryMappedBufferClosed() | ||
| Log.rebuildProducerState(producerStateManager, segments, logStartOffset, lastOffset, recordVersion, time, | ||
| reloadFromCleanShutdown = false) | ||
| reloadFromCleanShutdown = false, logIdent) | ||
| } | ||
|
|
||
| def activeProducers: Seq[DescribeProducersResponseData.ProducerState] = { | ||
|
|
@@ -1888,14 +1888,14 @@ class Log(@volatile private var _dir: File, | |
|
|
||
| private def deleteSegmentFiles(segments: Iterable[LogSegment], asyncDelete: Boolean, deleteProducerStateSnapshots: Boolean = true): Unit = { | ||
| Log.deleteSegmentFiles(segments, asyncDelete, deleteProducerStateSnapshots, dir, topicPartition, | ||
| config, scheduler, logDirFailureChannel, producerStateManager) | ||
| config, scheduler, logDirFailureChannel, producerStateManager, this.logIdent) | ||
| } | ||
|
|
||
| private[log] def replaceSegments(newSegments: Seq[LogSegment], oldSegments: Seq[LogSegment], isRecoveredSwapFile: Boolean = false): Unit = { | ||
| lock synchronized { | ||
| checkIfMemoryMappedBufferClosed() | ||
| Log.replaceSegments(segments, newSegments, oldSegments, isRecoveredSwapFile, dir, topicPartition, | ||
| config, scheduler, logDirFailureChannel, producerStateManager) | ||
| config, scheduler, logDirFailureChannel, producerStateManager, this.logIdent) | ||
| } | ||
| } | ||
|
|
||
|
|
@@ -1937,7 +1937,7 @@ class Log(@volatile private var _dir: File, | |
| } | ||
|
|
||
| private[log] def splitOverflowedSegment(segment: LogSegment): List[LogSegment] = lock synchronized { | ||
| Log.splitOverflowedSegment(segment, segments, dir, topicPartition, config, scheduler, logDirFailureChannel, producerStateManager) | ||
| Log.splitOverflowedSegment(segment, segments, dir, topicPartition, config, scheduler, logDirFailureChannel, producerStateManager, this.logIdent) | ||
| } | ||
|
|
||
| } | ||
|
|
@@ -2005,7 +2005,12 @@ object Log extends Logging { | |
| Files.createDirectories(dir.toPath) | ||
| val topicPartition = Log.parseTopicPartitionName(dir) | ||
| val segments = new LogSegments(topicPartition) | ||
| val leaderEpochCache = Log.maybeCreateLeaderEpochCache(dir, topicPartition, logDirFailureChannel, config.messageFormatVersion.recordVersion) | ||
| val leaderEpochCache = Log.maybeCreateLeaderEpochCache( | ||
| dir, | ||
| topicPartition, | ||
| logDirFailureChannel, | ||
| config.messageFormatVersion.recordVersion, | ||
| s"[Log partition=$topicPartition, dir=${dir.getParent}] )") | ||
| val producerStateManager = new ProducerStateManager(topicPartition, dir, maxProducerIdExpirationMs) | ||
| val offsets = LogLoader.load(LoadLogParams( | ||
| dir, | ||
|
|
@@ -2226,12 +2231,14 @@ object Log extends Logging { | |
| * @param topicPartition The topic partition | ||
| * @param logDirFailureChannel The LogDirFailureChannel to asynchronously handle log dir failure | ||
| * @param recordVersion The record version | ||
| * @param logPrefix The logging prefix | ||
| * @return The new LeaderEpochFileCache instance (if created), none otherwise | ||
| */ | ||
| def maybeCreateLeaderEpochCache(dir: File, | ||
| topicPartition: TopicPartition, | ||
| logDirFailureChannel: LogDirFailureChannel, | ||
| recordVersion: RecordVersion): Option[LeaderEpochFileCache] = { | ||
| recordVersion: RecordVersion, | ||
| logPrefix: String): Option[LeaderEpochFileCache] = { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Could we add the default
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Introducing a default value can lead to programming error, because we could forget to pass it when it is really needed to be passed.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Ok, good to me. |
||
| val leaderEpochFile = LeaderEpochCheckpointFile.newFile(dir) | ||
|
|
||
| def newLeaderEpochFileCache(): LeaderEpochFileCache = { | ||
|
|
@@ -2246,7 +2253,7 @@ object Log extends Logging { | |
| None | ||
|
|
||
| if (currentCache.exists(_.nonEmpty)) | ||
| warn(s"Deleting non-empty leader epoch cache due to incompatible message format $recordVersion") | ||
| warn(s"${logPrefix}Deleting non-empty leader epoch cache due to incompatible message format $recordVersion") | ||
|
|
||
| Files.deleteIfExists(leaderEpochFile.toPath) | ||
| None | ||
|
|
@@ -2293,6 +2300,7 @@ object Log extends Logging { | |
| * @param logDirFailureChannel The LogDirFailureChannel to asynchronously handle log dir failure | ||
| * @param producerStateManager The ProducerStateManager instance (if any) containing state associated | ||
| * with the existingSegments | ||
| * @param logPrefix The logging prefix | ||
| */ | ||
| private[log] def replaceSegments(existingSegments: LogSegments, | ||
| newSegments: Seq[LogSegment], | ||
|
|
@@ -2303,7 +2311,8 @@ object Log extends Logging { | |
| config: LogConfig, | ||
| scheduler: Scheduler, | ||
| logDirFailureChannel: LogDirFailureChannel, | ||
| producerStateManager: ProducerStateManager): Unit = { | ||
| producerStateManager: ProducerStateManager, | ||
| logPrefix: String): Unit = { | ||
| val sortedNewSegments = newSegments.sortBy(_.baseOffset) | ||
| // Some old segments may have been removed from index and scheduled for async deletion after the caller reads segments | ||
| // but before this method is executed. We want to filter out those segments to avoid calling asyncDeleteSegment() | ||
|
|
@@ -2332,7 +2341,8 @@ object Log extends Logging { | |
| config, | ||
| scheduler, | ||
| logDirFailureChannel, | ||
| producerStateManager) | ||
| producerStateManager, | ||
| logPrefix) | ||
|
kowshik marked this conversation as resolved.
|
||
| } | ||
| // okay we are safe now, remove the swap suffix | ||
| sortedNewSegments.foreach(_.changeFileSuffixes(Log.SwapFileSuffix, "")) | ||
|
|
@@ -2359,7 +2369,7 @@ object Log extends Logging { | |
| * @param logDirFailureChannel The LogDirFailureChannel to asynchronously handle log dir failure | ||
| * @param producerStateManager The ProducerStateManager instance (if any) containing state associated | ||
| * with the existingSegments | ||
| * | ||
| * @param logPrefix The logging prefix | ||
| * @throws IOException if the file can't be renamed and still exists | ||
| */ | ||
| private[log] def deleteSegmentFiles(segmentsToDelete: Iterable[LogSegment], | ||
|
|
@@ -2370,11 +2380,12 @@ object Log extends Logging { | |
| config: LogConfig, | ||
| scheduler: Scheduler, | ||
| logDirFailureChannel: LogDirFailureChannel, | ||
| producerStateManager: ProducerStateManager): Unit = { | ||
| producerStateManager: ProducerStateManager, | ||
| logPrefix: String): Unit = { | ||
| segmentsToDelete.foreach(_.changeFileSuffixes("", Log.DeletedFileSuffix)) | ||
|
|
||
| def deleteSegments(): Unit = { | ||
| info(s"Deleting segment files ${segmentsToDelete.mkString(",")}") | ||
| info(s"${logPrefix}Deleting segment files ${segmentsToDelete.mkString(",")}") | ||
| val parentDir = dir.getParent | ||
| maybeHandleIOException(logDirFailureChannel, parentDir, s"Error while deleting segments for $topicPartition in dir $parentDir") { | ||
| segmentsToDelete.foreach { segment => | ||
|
|
@@ -2429,14 +2440,16 @@ object Log extends Logging { | |
| * @param time The time instance used for checking the clock | ||
| * @param reloadFromCleanShutdown True if the producer state is being built after a clean shutdown, | ||
| * false otherwise. | ||
| * @param logPrefix The logging prefix | ||
| */ | ||
| private[log] def rebuildProducerState(producerStateManager: ProducerStateManager, | ||
| segments: LogSegments, | ||
| logStartOffset: Long, | ||
| lastOffset: Long, | ||
| recordVersion: RecordVersion, | ||
| time: Time, | ||
| reloadFromCleanShutdown: Boolean): Unit = { | ||
| reloadFromCleanShutdown: Boolean, | ||
| logPrefix: String): Unit = { | ||
| val allSegments = segments.values | ||
| val offsetsToSnapshot = | ||
| if (allSegments.nonEmpty) { | ||
|
|
@@ -2445,7 +2458,7 @@ object Log extends Logging { | |
| } else { | ||
| Seq(Some(lastOffset)) | ||
| } | ||
| info(s"Loading producer state till offset $lastOffset with message format version ${recordVersion.value}") | ||
| info(s"${logPrefix}Loading producer state till offset $lastOffset with message format version ${recordVersion.value}") | ||
|
|
||
| // We want to avoid unnecessary scanning of the log to build the producer state when the broker is being | ||
| // upgraded. The basic idea is to use the absence of producer snapshot files to detect the upgrade case, | ||
|
|
@@ -2469,7 +2482,7 @@ object Log extends Logging { | |
| producerStateManager.takeSnapshot() | ||
| } | ||
| } else { | ||
| info(s"Reloading from producer snapshot and rebuilding producer state from offset $lastOffset") | ||
| info(s"${logPrefix}Reloading from producer snapshot and rebuilding producer state from offset $lastOffset") | ||
| val isEmptyBeforeTruncation = producerStateManager.isEmpty && producerStateManager.mapEndOffset >= lastOffset | ||
| val producerStateLoadStart = time.milliseconds() | ||
| producerStateManager.truncateAndReload(logStartOffset, lastOffset, time.milliseconds()) | ||
|
|
@@ -2508,7 +2521,7 @@ object Log extends Logging { | |
| } | ||
| producerStateManager.updateMapEndOffset(lastOffset) | ||
| producerStateManager.takeSnapshot() | ||
| info(s"Producer state recovery took ${segmentRecoveryStart - producerStateLoadStart}ms for snapshot load " + | ||
| info(s"${logPrefix}Producer state recovery took ${segmentRecoveryStart - producerStateLoadStart}ms for snapshot load " + | ||
| s"and ${time.milliseconds() - segmentRecoveryStart}ms for segment recovery from offset $lastOffset") | ||
| } | ||
| } | ||
|
|
@@ -2535,6 +2548,7 @@ object Log extends Logging { | |
| * @param logDirFailureChannel The LogDirFailureChannel to asynchronously handle log dir failure | ||
| * @param producerStateManager The ProducerStateManager instance (if any) containing state associated | ||
| * with the existingSegments | ||
| * @param logPrefix The logging prefix | ||
| * @return List of new segments that replace the input segment | ||
| */ | ||
| private[log] def splitOverflowedSegment(segment: LogSegment, | ||
|
|
@@ -2544,11 +2558,12 @@ object Log extends Logging { | |
| config: LogConfig, | ||
| scheduler: Scheduler, | ||
| logDirFailureChannel: LogDirFailureChannel, | ||
| producerStateManager: ProducerStateManager): List[LogSegment] = { | ||
| producerStateManager: ProducerStateManager, | ||
| logPrefix: String): List[LogSegment] = { | ||
| require(Log.isLogFile(segment.log.file), s"Cannot split file ${segment.log.file.getAbsoluteFile}") | ||
| require(segment.hasOverflow, "Split operation is only permitted for segments with overflow") | ||
|
|
||
| info(s"Splitting overflowed segment $segment") | ||
| info(s"${logPrefix}Splitting overflowed segment $segment") | ||
|
|
||
| val newSegments = ListBuffer[LogSegment]() | ||
| try { | ||
|
|
@@ -2581,9 +2596,9 @@ object Log extends Logging { | |
| s" before: ${segment.log.sizeInBytes} after: $totalSizeOfNewSegments") | ||
|
|
||
| // replace old segment with new ones | ||
| info(s"Replacing overflowed segment $segment with split segments $newSegments") | ||
| info(s"${logPrefix}Replacing overflowed segment $segment with split segments $newSegments") | ||
| replaceSegments(existingSegments, newSegments.toList, List(segment), isRecoveredSwapFile = false, | ||
| dir, topicPartition, config, scheduler, logDirFailureChannel, producerStateManager) | ||
| dir, topicPartition, config, scheduler, logDirFailureChannel, producerStateManager, logPrefix) | ||
| newSegments.toList | ||
| } catch { | ||
| case e: Exception => | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.