Skip to content
Closed
Show file tree
Hide file tree
Changes from 11 commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
9fd69d7
Add HostedPartition.Fenced to ReplicaManager
rondagostino Jan 22, 2021
fef34f6
implement makeFenced()
rondagostino Jan 23, 2021
fdc1f5b
cleanups
rondagostino Jan 23, 2021
d9fa4dd
Implement unfencing of leader partitions
rondagostino Jan 25, 2021
0324d5b
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 25, 2021
f4bf859
Implemented unfenceFencedFollowerPartitions()
rondagostino Jan 25, 2021
2a78717
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 25, 2021
d667f38
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 26, 2021
776eae4
Eliminate fencedPartitionStates
rondagostino Jan 26, 2021
8620ea4
Eliminate "get" prefix on long method names
rondagostino Jan 26, 2021
fe98cd6
Rename "Fenced" to "Deferred"
rondagostino Jan 26, 2021
d2d688b
renames for clarity/risk mitigation, respond to some review comments
rondagostino Jan 27, 2021
1f019dc
Defer changes when using Raft only/fix failng tests
rondagostino Jan 27, 2021
c4487d5
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 27, 2021
a4d238b
invoke onLeadershipChange() when applying deferred changes
rondagostino Jan 27, 2021
09c5966
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 27, 2021
eebc368
Reuse makeLeader/makeFollower helpers
rondagostino Jan 27, 2021
993ca24
Do not expose deferred partitions
rondagostino Jan 27, 2021
275ef23
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 28, 2021
32bcfec
Merge remote-tracking branch 'origin/kip-500' into rtd_HostedPartitio…
rondagostino Jan 28, 2021
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
4 changes: 3 additions & 1 deletion core/src/main/scala/kafka/log/LogManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -346,7 +346,9 @@ class LogManager(logDirs: Seq[File],
val numLogsLoaded = new AtomicInteger(0)
numTotalLogs += logsToLoad.length

val jobsForDir = logsToLoad.map { logDir =>
val jobsForDir = logsToLoad
.filter(logDir => Log.parseTopicPartitionName(logDir).topic != "@metadata")
.map { logDir =>
val runnable: Runnable = () => {
try {
debug(s"Loading log $logDir")
Expand Down
2 changes: 2 additions & 0 deletions core/src/main/scala/kafka/server/DelayedDeleteRecords.scala
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,8 @@ class DelayedDeleteRecords(delayMs: Long,
case None =>
(false, Errors.NOT_LEADER_OR_FOLLOWER, DeleteRecordsResponse.INVALID_LOW_WATERMARK)
}
case HostedPartition.Deferred(_) =>
(false, Errors.NOT_LEADER_OR_FOLLOWER, DeleteRecordsResponse.INVALID_LOW_WATERMARK)

case HostedPartition.Offline =>
(false, Errors.KAFKA_STORAGE_ERROR, DeleteRecordsResponse.INVALID_LOW_WATERMARK)
Expand Down
4 changes: 2 additions & 2 deletions core/src/main/scala/kafka/server/DelayedFetch.scala
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ class DelayedFetch(delayMs: Long,
val fetchLeaderEpoch = fetchStatus.fetchInfo.currentLeaderEpoch
try {
if (fetchOffset != LogOffsetMetadata.UnknownOffsetMetadata) {
val partition = replicaManager.getPartitionOrException(topicPartition)
val partition = replicaManager.onlinePartitionOrException(topicPartition)
val offsetSnapshot = partition.fetchOffsetSnapshot(fetchLeaderEpoch, fetchMetadata.fetchOnlyLeader)

val endOffset = fetchMetadata.fetchIsolation match {
Expand Down Expand Up @@ -184,7 +184,7 @@ class DelayedFetch(delayMs: Long,

val fetchPartitionData = logReadResults.map { case (tp, result) =>
val isReassignmentFetch = fetchMetadata.isFromFollower &&
replicaManager.isAddingReplica(tp, fetchMetadata.replicaId)
replicaManager.isAddingReplica(tp, fetchMetadata.replicaId, false)

tp -> FetchPartitionData(
result.error,
Expand Down
2 changes: 1 addition & 1 deletion core/src/main/scala/kafka/server/DelayedProduce.scala
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ class DelayedProduce(delayMs: Long,
trace(s"Checking produce satisfaction for $topicPartition, current status $status")
// skip those partitions that have already been satisfied
if (status.acksPending) {
val (hasEnough, error) = replicaManager.getPartitionOrError(topicPartition) match {
val (hasEnough, error) = replicaManager.onlinePartitionOrError(topicPartition) match {
case Left(err) =>
// Case A
(false, err)
Expand Down
10 changes: 5 additions & 5 deletions core/src/main/scala/kafka/server/ReplicaAlterLogDirsThread.scala
Original file line number Diff line number Diff line change
Expand Up @@ -148,12 +148,12 @@ class ReplicaAlterLogDirsThread(name: String,
}

override protected def fetchEarliestOffsetFromLeader(topicPartition: TopicPartition, leaderEpoch: Int): Long = {
val partition = replicaMgr.getPartitionOrException(topicPartition)
val partition = replicaMgr.onlinePartitionOrException(topicPartition)
partition.localLogOrException.logStartOffset
}

override protected def fetchLatestOffsetFromLeader(topicPartition: TopicPartition, leaderEpoch: Int): Long = {
val partition = replicaMgr.getPartitionOrException(topicPartition)
val partition = replicaMgr.onlinePartitionOrException(topicPartition)
partition.localLogOrException.logEndOffset
}

Expand All @@ -170,7 +170,7 @@ class ReplicaAlterLogDirsThread(name: String,
.setPartition(tp.partition)
.setErrorCode(Errors.NONE.code)
} else {
val partition = replicaMgr.getPartitionOrException(tp)
val partition = replicaMgr.onlinePartitionOrException(tp)
partition.lastOffsetForLeaderEpoch(
currentLeaderEpoch = epochData.currentLeaderEpoch,
leaderEpoch = epochData.leaderEpoch,
Expand Down Expand Up @@ -206,12 +206,12 @@ class ReplicaAlterLogDirsThread(name: String,
* exchange with the current replica to truncate to the largest common log prefix for the topic partition
*/
override def truncate(topicPartition: TopicPartition, truncationState: OffsetTruncationState): Unit = {
val partition = replicaMgr.getPartitionOrException(topicPartition)
val partition = replicaMgr.onlinePartitionOrException(topicPartition)
partition.truncateTo(truncationState.offset, isFuture = true)
}

override protected def truncateFullyAndStartAt(topicPartition: TopicPartition, offset: Long): Unit = {
val partition = replicaMgr.getPartitionOrException(topicPartition)
val partition = replicaMgr.onlinePartitionOrException(topicPartition)
partition.truncateFullyAndStartAt(offset, isFuture = true)
}

Expand Down
8 changes: 4 additions & 4 deletions core/src/main/scala/kafka/server/ReplicaFetcherThread.scala
Original file line number Diff line number Diff line change
Expand Up @@ -109,19 +109,19 @@ class ReplicaFetcherThread(name: String,
val fetchSessionHandler = new FetchSessionHandler(logContext, sourceBroker.id)

override protected def latestEpoch(topicPartition: TopicPartition): Option[Int] = {
replicaMgr.localLogOrException(topicPartition).latestEpoch
replicaMgr.localOnlineLogOrException(topicPartition).latestEpoch
}

override protected def logStartOffset(topicPartition: TopicPartition): Long = {
replicaMgr.localLogOrException(topicPartition).logStartOffset
replicaMgr.localOnlineLogOrException(topicPartition).logStartOffset
}

override protected def logEndOffset(topicPartition: TopicPartition): Long = {
replicaMgr.localLogOrException(topicPartition).logEndOffset
replicaMgr.localOnlineLogOrException(topicPartition).logEndOffset
}

override protected def endOffsetForEpoch(topicPartition: TopicPartition, epoch: Int): Option[OffsetAndEpoch] = {
replicaMgr.localLogOrException(topicPartition).endOffsetForEpoch(epoch)
replicaMgr.localOnlineLogOrException(topicPartition).endOffsetForEpoch(epoch)
}

override def initiateShutdown(): Boolean = {
Expand Down
Loading