Skip to content
Merged
Show file tree
Hide file tree
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
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,8 @@
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.function.LongSupplier;
import java.util.function.Predicate;
import java.util.function.Supplier;
import java.util.regex.Pattern;
import java.util.stream.Collector;
import java.util.stream.Collectors;
Expand Down Expand Up @@ -516,7 +516,7 @@ synchronized void updateLastStableOffset(TopicPartition tp, long lastStableOffse
* @param preferredReadReplicaId The preferred read replica
* @param timeMs The time at which this preferred replica is no longer valid
*/
public synchronized void updatePreferredReadReplica(TopicPartition tp, int preferredReadReplicaId, Supplier<Long> timeMs) {
public synchronized void updatePreferredReadReplica(TopicPartition tp, int preferredReadReplicaId, LongSupplier timeMs) {
assignedState(tp).updatePreferredReadReplica(preferredReadReplicaId, timeMs);
}

Expand Down Expand Up @@ -721,10 +721,10 @@ private Optional<Integer> preferredReadReplica(long timeMs) {
}
}

private void updatePreferredReadReplica(int preferredReadReplica, Supplier<Long> timeMs) {
private void updatePreferredReadReplica(int preferredReadReplica, LongSupplier timeMs) {
if (this.preferredReadReplica == null || preferredReadReplica != this.preferredReadReplica) {
this.preferredReadReplica = preferredReadReplica;
this.preferredReadReplicaExpireTimeMs = timeMs.get();
this.preferredReadReplicaExpireTimeMs = timeMs.getAsLong();
}
}

Expand Down
57 changes: 24 additions & 33 deletions core/src/main/scala/kafka/server/ReplicaManager.scala
Original file line number Diff line number Diff line change
Expand Up @@ -1053,7 +1053,7 @@ class ReplicaManager(val config: KafkaConfig,
metadata => findPreferredReadReplica(partition, metadata, replicaId, fetchInfo.fetchOffset, fetchTimeMs))

if (preferredReadReplica.isDefined) {
replicaSelectorOpt.foreach{ selector =>
replicaSelectorOpt.foreach { selector =>
debug(s"Replica selector ${selector.getClass.getSimpleName} returned preferred replica " +
s"${preferredReadReplica.get} for $clientMetadata")
}
Expand All @@ -1079,9 +1079,9 @@ class ReplicaManager(val config: KafkaConfig,
fetchOnlyFromLeader = fetchOnlyFromLeader,
minOneMessage = minOneMessage)

// Check if the HW known to the follower is behind the actual HW
val followerNeedsHwUpdate: Boolean = partition.getReplica(replicaId)
.exists(replica => replica.lastSentHighWatermark < readInfo.highWatermark)
// Check if the HW known to the follower is behind the actual HW if a replica selector is defined
val followerNeedsHwUpdate = replicaSelectorOpt.isDefined &&
partition.getReplica(replicaId).exists(replica => replica.lastSentHighWatermark < readInfo.highWatermark)

val fetchDataInfo = if (shouldLeaderThrottle(quota, partition, replicaId)) {
// If the partition is being throttled, simply return an empty set.
Expand Down Expand Up @@ -1170,44 +1170,35 @@ class ReplicaManager(val config: KafkaConfig,
replicaId: Int,
fetchOffset: Long,
currentTimeMs: Long): Option[Int] = {
if (partition.isLeader) {
if (Request.isValidBrokerId(replicaId)) {
// Don't look up preferred for follower fetches via normal replication
Option.empty
} else {
partition.leaderReplicaIdOpt.flatMap { leaderReplicaId =>
// Don't look up preferred for follower fetches via normal replication
if (Request.isValidBrokerId(replicaId))
None
else {
replicaSelectorOpt.flatMap { replicaSelector =>
val replicaEndpoints = metadataCache.getPartitionReplicaEndpoints(partition.topicPartition, new ListenerName(clientMetadata.listenerName))
var replicaInfoSet: Set[ReplicaView] = partition.remoteReplicas
val replicaEndpoints = metadataCache.getPartitionReplicaEndpoints(partition.topicPartition,
new ListenerName(clientMetadata.listenerName))
val replicaInfos = partition.remoteReplicas
// Exclude replicas that don't have the requested offset (whether or not if they're in the ISR)
.filter(replica => replica.logEndOffset >= fetchOffset)
.filter(replica => replica.logStartOffset <= fetchOffset)
.filter(replica => replica.logEndOffset >= fetchOffset && replica.logStartOffset <= fetchOffset)
.map(replica => new DefaultReplicaView(
replicaEndpoints.getOrElse(replica.brokerId, Node.noNode()),
replica.logEndOffset,
currentTimeMs - replica.lastCaughtUpTimeMs))
.toSet

if (partition.leaderReplicaIdOpt.isDefined) {
val leaderReplica: ReplicaView = partition.leaderReplicaIdOpt
.map(replicaId => replicaEndpoints.getOrElse(replicaId, Node.noNode()))
.map(leaderNode => new DefaultReplicaView(leaderNode, partition.localLogOrException.logEndOffset, 0L))
.get
replicaInfoSet ++= Set(leaderReplica)

val partitionInfo = new DefaultPartitionView(replicaInfoSet.asJava, leaderReplica)
replicaSelector.select(partition.topicPartition, clientMetadata, partitionInfo).asScala
.filter(!_.endpoint.isEmpty)
// Even though the replica selector can return the leader, we don't want to send it out with the
// FetchResponse, so we exclude it here
.filter(!_.equals(leaderReplica))
.map(_.endpoint.id)
} else {
None

val leaderReplica = new DefaultReplicaView(
replicaEndpoints.getOrElse(leaderReplicaId, Node.noNode()),
partition.localLogOrException.logEndOffset, 0L)
val replicaInfoSet = mutable.Set[ReplicaView]() ++= replicaInfos += leaderReplica

val partitionInfo = new DefaultPartitionView(replicaInfoSet.asJava, leaderReplica)
replicaSelector.select(partition.topicPartition, clientMetadata, partitionInfo).asScala.collect {
// Even though the replica selector can return the leader, we don't want to send it out with the
// FetchResponse, so we exclude it here
case selected if !selected.endpoint.isEmpty && selected != leaderReplica => selected.endpoint.id
}
}
}
} else {
None
}
}

Expand Down
Loading