Skip to content
Merged
Changes from all 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
57 changes: 30 additions & 27 deletions core/src/main/scala/kafka/log/Log.scala
Original file line number Diff line number Diff line change
Expand Up @@ -935,7 +935,7 @@ class Log(@volatile private var _dir: File,
maxPosition = maxPosition,
minOneMessage = false)
if (fetchDataInfo != null)
loadProducersFromLog(producerStateManager, fetchDataInfo.records)
loadProducersFromRecords(producerStateManager, fetchDataInfo.records)
}
}
producerStateManager.updateMapEndOffset(lastOffset)
Expand All @@ -948,22 +948,6 @@ class Log(@volatile private var _dir: File,
maybeIncrementFirstUnstableOffset()
}

private def loadProducersFromLog(producerStateManager: ProducerStateManager, records: Records): Unit = {
val loadedProducers = mutable.Map.empty[Long, ProducerAppendInfo]
val completedTxns = ListBuffer.empty[CompletedTxn]
records.batches.forEach { batch =>
if (batch.hasProducerId) {
val maybeCompletedTxn = updateProducers(batch,
loadedProducers,
firstOffsetMetadata = None,
origin = AppendOrigin.Replication)
maybeCompletedTxn.foreach(completedTxns += _)
}
}
loadedProducers.values.foreach(producerStateManager.update)
completedTxns.foreach(producerStateManager.completeTxn)
}

private[log] def activeProducersWithLastSequence: Map[Long, Int] = lock synchronized {
producerStateManager.activeProducers.map { case (producerId, producerIdEntry) =>
(producerId, producerIdEntry.lastSeq)
Expand Down Expand Up @@ -1349,7 +1333,7 @@ class Log(@volatile private var _dir: File,
else
None

val maybeCompletedTxn = updateProducers(batch, updatedProducers, firstOffsetMetadata, origin)
val maybeCompletedTxn = updateProducers(producerStateManager, batch, updatedProducers, firstOffsetMetadata, origin)
maybeCompletedTxn.foreach(completedTxns += _)
}

Expand Down Expand Up @@ -1456,15 +1440,6 @@ class Log(@volatile private var _dir: File,
RecordConversionStats.EMPTY, sourceCodec, targetCodec, shallowMessageCount, validBytesCount, monotonic, lastOffsetOfFirstBatch)
}

private def updateProducers(batch: RecordBatch,
producers: mutable.Map[Long, ProducerAppendInfo],
firstOffsetMetadata: Option[LogOffsetMetadata],
origin: AppendOrigin): Option[CompletedTxn] = {
val producerId = batch.producerId
val appendInfo = producers.getOrElseUpdate(producerId, producerStateManager.prepareUpdate(producerId, origin))
appendInfo.append(batch, firstOffsetMetadata)
}

/**
* Trim any invalid bytes from the end of this message set (if there are any)
*
Expand Down Expand Up @@ -2670,6 +2645,34 @@ object Log {
private def isLogFile(file: File): Boolean =
file.getPath.endsWith(LogFileSuffix)

private def loadProducersFromRecords(producerStateManager: ProducerStateManager, records: Records): Unit = {
val loadedProducers = mutable.Map.empty[Long, ProducerAppendInfo]
val completedTxns = ListBuffer.empty[CompletedTxn]
records.batches.forEach { batch =>
if (batch.hasProducerId) {
val maybeCompletedTxn = updateProducers(
producerStateManager,
batch,
loadedProducers,
firstOffsetMetadata = None,
origin = AppendOrigin.Replication)
maybeCompletedTxn.foreach(completedTxns += _)
}
}
loadedProducers.values.foreach(producerStateManager.update)
completedTxns.foreach(producerStateManager.completeTxn)
}

private def updateProducers(producerStateManager: ProducerStateManager,
batch: RecordBatch,
producers: mutable.Map[Long, ProducerAppendInfo],
firstOffsetMetadata: Option[LogOffsetMetadata],
origin: AppendOrigin): Option[CompletedTxn] = {
val producerId = batch.producerId
val appendInfo = producers.getOrElseUpdate(producerId, producerStateManager.prepareUpdate(producerId, origin))
appendInfo.append(batch, firstOffsetMetadata)
}

}

object LogMetricNames {
Expand Down