Skip to content
24 changes: 14 additions & 10 deletions clients/src/main/java/org/apache/kafka/clients/Metadata.java
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,8 @@ public synchronized int requestUpdate() {
*/
public synchronized boolean updateLastSeenEpochIfNewer(TopicPartition topicPartition, int leaderEpoch) {
Objects.requireNonNull(topicPartition, "TopicPartition cannot be null");
if (leaderEpoch < 0)
throw new IllegalArgumentException("Invalid leader epoch " + leaderEpoch + " (must be non-negative)");
return updateLastSeenEpoch(topicPartition, leaderEpoch, oldEpoch -> leaderEpoch > oldEpoch, true);
}

Expand Down Expand Up @@ -211,19 +213,15 @@ synchronized Optional<MetadataResponse.PartitionMetadata> partitionMetadataIfCur
}
}

public synchronized Optional<LeaderAndEpoch> currentLeader(TopicPartition topicPartition) {
public synchronized LeaderAndEpoch currentLeader(TopicPartition topicPartition) {
Optional<MetadataResponse.PartitionMetadata> maybeMetadata = partitionMetadataIfCurrent(topicPartition);
if (!maybeMetadata.isPresent())
return Optional.empty();
return new LeaderAndEpoch(Optional.empty(), Optional.ofNullable(lastSeenLeaderEpochs.get(topicPartition)));

MetadataResponse.PartitionMetadata partitionMetadata = maybeMetadata.get();
Optional<Integer> leaderEpoch = partitionMetadata.leaderEpoch;
Comment thread
ijuma marked this conversation as resolved.
Outdated
Optional<Node> leaderNodeOpt = cache.nodeById(partitionMetadata.leaderId);
return leaderNodeOpt.map(leaderNode -> new LeaderAndEpoch(leaderNode, leaderEpoch));
}

public synchronized LeaderAndEpoch currentLeaderOrEmpty(TopicPartition tp) {
return currentLeader(tp).orElse(new LeaderAndEpoch(Node.noNode(), lastSeenLeaderEpoch(tp)));
return new LeaderAndEpoch(leaderNodeOpt, leaderEpoch);
}

public synchronized void bootstrap(List<InetSocketAddress> addresses) {
Expand Down Expand Up @@ -508,14 +506,20 @@ private MetadataRequestAndVersion(MetadataRequest.Builder requestBuilder,
}
}

/**
* Represents current leader state known in metadata. It is possible that we know the leader, but not the
* epoch if the metadata is received from a broker which does not support a sufficient Metadata API version.
* It is also possible that we know of the leader epoch, but not the leader when it is derived
* from an external source (e.g. a committed offset).
*/
public static class LeaderAndEpoch {

public static final LeaderAndEpoch NO_LEADER_OR_EPOCH = new LeaderAndEpoch(Node.noNode(), Optional.empty());
public static final LeaderAndEpoch NO_LEADER_OR_EPOCH = new LeaderAndEpoch(Optional.empty(), Optional.empty());
Comment thread
ijuma marked this conversation as resolved.
Outdated

public final Node leader;
public final Optional<Node> leader;
public final Optional<Integer> epoch;

public LeaderAndEpoch(Node leader, Optional<Integer> epoch) {
public LeaderAndEpoch(Optional<Node> leader, Optional<Integer> epoch) {
this.leader = Objects.requireNonNull(leader);
this.epoch = Objects.requireNonNull(epoch);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1604,7 +1604,7 @@ public void seek(TopicPartition partition, long offset) {
SubscriptionState.FetchPosition newPosition = new SubscriptionState.FetchPosition(
offset,
Optional.empty(), // This will ensure we skip validation
this.metadata.currentLeaderOrEmpty(partition));
this.metadata.currentLeader(partition));
this.subscriptions.seekUnvalidated(partition, newPosition);
} finally {
release();
Expand Down Expand Up @@ -1635,7 +1635,7 @@ public void seek(TopicPartition partition, OffsetAndMetadata offsetAndMetadata)
} else {
log.info("Seeking to offset {} for partition {}", offset, partition);
}
Metadata.LeaderAndEpoch currentLeaderAndEpoch = this.metadata.currentLeaderOrEmpty(partition);
Metadata.LeaderAndEpoch currentLeaderAndEpoch = this.metadata.currentLeader(partition);
SubscriptionState.FetchPosition newPosition = new SubscriptionState.FetchPosition(
offsetAndMetadata.offset(),
offsetAndMetadata.leaderEpoch(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.Metric;
import org.apache.kafka.common.MetricName;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.WakeupException;
Expand All @@ -37,6 +36,7 @@
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Queue;
import java.util.Set;
import java.util.concurrent.TimeUnit;
Expand Down Expand Up @@ -202,8 +202,9 @@ public synchronized ConsumerRecords<K, V> poll(final Duration timeout) {

if (assignment().contains(entry.getKey()) && rec.offset() >= position) {
results.computeIfAbsent(entry.getKey(), partition -> new ArrayList<>()).add(rec);
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(Optional.empty(), rec.leaderEpoch());
SubscriptionState.FetchPosition newPosition = new SubscriptionState.FetchPosition(
rec.offset() + 1, rec.leaderEpoch(), new Metadata.LeaderAndEpoch(Node.noNode(), rec.leaderEpoch()));
rec.offset() + 1, rec.leaderEpoch(), leaderAndEpoch);
subscriptions.position(entry.getKey(), newPosition);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package org.apache.kafka.clients.consumer;

import org.apache.kafka.common.requests.OffsetFetchResponse;
import org.apache.kafka.common.requests.RequestUtils;

import java.io.Serializable;
import java.util.Objects;
Expand Down Expand Up @@ -94,7 +95,9 @@ public String metadata() {
* @return the leader epoch or empty if not known
*/
public Optional<Integer> leaderEpoch() {
return Optional.ofNullable(leaderEpoch);
if (leaderEpoch == null || leaderEpoch < 0)
return Optional.empty();
return Optional.of(leaderEpoch);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -766,7 +766,7 @@ public boolean refreshCommittedOffsetsIfNeeded(Timer timer) {
// it's possible that the partition is no longer assigned when the response is received,
// so we need to ignore seeking if that's the case
if (this.subscriptions.isAssigned(tp)) {
final ConsumerMetadata.LeaderAndEpoch leaderAndEpoch = metadata.currentLeaderOrEmpty(tp);
final ConsumerMetadata.LeaderAndEpoch leaderAndEpoch = metadata.currentLeader(tp);
final SubscriptionState.FetchPosition position = new SubscriptionState.FetchPosition(
offsetAndMetadata.offset(), offsetAndMetadata.leaderEpoch(),
leaderAndEpoch);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -485,7 +485,7 @@ public void validateOffsetsIfNeeded() {

// Validate each partition against the current leader and epoch
subscriptions.assignedPartitions().forEach(topicPartition -> {
ConsumerMetadata.LeaderAndEpoch leaderAndEpoch = metadata.currentLeaderOrEmpty(topicPartition);
ConsumerMetadata.LeaderAndEpoch leaderAndEpoch = metadata.currentLeader(topicPartition);
subscriptions.maybeValidatePositionForCurrentLeader(topicPartition, leaderAndEpoch);
});

Expand Down Expand Up @@ -716,7 +716,7 @@ private List<ConsumerRecord<K, V>> fetchRecords(CompletedFetch completedFetch, i

private void resetOffsetIfNeeded(TopicPartition partition, OffsetResetStrategy requestedResetStrategy, ListOffsetData offsetData) {
SubscriptionState.FetchPosition position = new SubscriptionState.FetchPosition(
offsetData.offset, offsetData.leaderEpoch, metadata.currentLeaderOrEmpty(partition));
offsetData.offset, offsetData.leaderEpoch, metadata.currentLeader(partition));
offsetData.leaderEpoch.ifPresent(epoch -> metadata.updateLastSeenEpochIfNewer(partition, epoch));
subscriptions.maybeSeekUnvalidated(partition, position.offset, requestedResetStrategy);
}
Expand Down Expand Up @@ -904,15 +904,14 @@ private Map<Node, Map<TopicPartition, ListOffsetRequest.PartitionData>> groupLis
for (Map.Entry<TopicPartition, Long> entry: timestampsToSearch.entrySet()) {
TopicPartition tp = entry.getKey();
Long offset = entry.getValue();
Optional<Metadata.LeaderAndEpoch> leaderAndEpochOpt = metadata.currentLeader(tp);
if (!leaderAndEpochOpt.isPresent()) {
Metadata.LeaderAndEpoch leaderAndEpoch = metadata.currentLeader(tp);

if (!leaderAndEpoch.leader.isPresent()) {
log.debug("Leader for partition {} is unknown for fetching offset {}", tp, offset);
metadata.requestUpdate();
partitionsToRetry.add(tp);
} else {
Metadata.LeaderAndEpoch leaderAndEpoch = leaderAndEpochOpt.get();
Node leader = leaderAndEpoch.leader;

Node leader = leaderAndEpoch.leader.get();
if (client.isUnavailable(leader)) {
client.maybeThrowAuthFailure(leader);

Expand Down Expand Up @@ -1100,18 +1099,21 @@ private Map<Node, FetchSessionHandler.FetchRequestData> prepareFetchRequests() {

// Ensure the position has an up-to-date leader
subscriptions.assignedPartitions().forEach(
tp -> subscriptions.maybeValidatePositionForCurrentLeader(tp, metadata.currentLeaderOrEmpty(tp)));
tp -> subscriptions.maybeValidatePositionForCurrentLeader(tp, metadata.currentLeader(tp)));

long currentTimeMs = time.milliseconds();

for (TopicPartition partition : fetchablePartitions()) {
// Use the preferred read replica if set, or the position's leader
SubscriptionState.FetchPosition position = this.subscriptions.position(partition);
Node node = selectReadReplica(partition, position.currentLeader.leader, currentTimeMs);

if (node == null || node.isEmpty()) {
Optional<Node> leaderOpt = position.currentLeader.leader;
if (!leaderOpt.isPresent()) {
metadata.requestUpdate();
} else if (client.isUnavailable(node)) {
continue;
}

Node node = selectReadReplica(partition, leaderOpt.get(), currentTimeMs);
if (client.isUnavailable(node)) {
client.maybeThrowAuthFailure(node);

// If we try to send during the reconnect blackout window, then the request is just
Expand Down Expand Up @@ -1153,7 +1155,8 @@ private Map<Node, Map<TopicPartition, SubscriptionState.FetchPosition>> regroupF
Map<TopicPartition, SubscriptionState.FetchPosition> partitionMap) {
return partitionMap.entrySet()
.stream()
.collect(Collectors.groupingBy(entry -> entry.getValue().currentLeader.leader,
.filter(entry -> entry.getValue().currentLeader.leader.isPresent())
.collect(Collectors.groupingBy(entry -> entry.getValue().currentLeader.leader.get(),
Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,20 +16,19 @@
*/
package org.apache.kafka.clients.consumer.internals;

import java.util.ArrayList;
import org.apache.kafka.clients.Metadata;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.NoOffsetForPartitionException;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.common.IsolationLevel;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.internals.PartitionStates;
import org.apache.kafka.common.requests.EpochEndOffset;
import org.apache.kafka.common.utils.LogContext;
import org.slf4j.Logger;

import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
Expand Down Expand Up @@ -985,7 +984,7 @@ public static class FetchPosition {
final Metadata.LeaderAndEpoch currentLeader;

FetchPosition(long offset) {
this(offset, Optional.empty(), new Metadata.LeaderAndEpoch(Node.noNode(), Optional.empty()));
this(offset, Optional.empty(), Metadata.LeaderAndEpoch.noLeaderOrEpoch());
}

public FetchPosition(long offset, Optional<Integer> offsetEpoch, Metadata.LeaderAndEpoch currentLeader) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ static Optional<Integer> getLeaderEpoch(Struct struct, Field.Int32 leaderEpochFi
return leaderEpochOpt;
}

static Optional<Integer> getLeaderEpoch(int leaderEpoch) {
public static Optional<Integer> getLeaderEpoch(int leaderEpoch) {
Optional<Integer> leaderEpochOpt = leaderEpoch == RecordBatch.NO_PARTITION_LEADER_EPOCH ?
Optional.empty() : Optional.of(leaderEpoch);
return leaderEpochOpt;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -82,4 +82,4 @@ public void testMissingLeaderEndpoint() {
assertEquals(nodesById.get(7), replicas.get(7));
}

}
}
Original file line number Diff line number Diff line change
Expand Up @@ -666,7 +666,7 @@ public void testLeaderMetadataInconsistentWithBrokerMetadata() {

assertNull(metadata.fetch().leaderFor(tp));
assertEquals(Optional.of(10), metadata.lastSeenLeaderEpoch(tp));
assertFalse(metadata.currentLeader(tp).isPresent());
assertFalse(metadata.currentLeader(tp).leader.isPresent());
}

private MetadataResponseTopicCollection buildTopicCollection(String topic, MetadataResponsePartition partitionMetadata) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3023,8 +3023,8 @@ public void testConsumingViaIncrementalFetchRequests() {

List<ConsumerRecord<byte[], byte[]>> records;
assignFromUser(new HashSet<>(Arrays.asList(tp0, tp1)));
subscriptions.seekValidated(tp0, new SubscriptionState.FetchPosition(0, Optional.empty(), metadata.currentLeaderOrEmpty(tp0)));
subscriptions.seekValidated(tp1, new SubscriptionState.FetchPosition(1, Optional.empty(), metadata.currentLeaderOrEmpty(tp1)));
subscriptions.seekValidated(tp0, new SubscriptionState.FetchPosition(0, Optional.empty(), metadata.currentLeader(tp0)));
subscriptions.seekValidated(tp1, new SubscriptionState.FetchPosition(1, Optional.empty(), metadata.currentLeader(tp1)));

// Fetch some records and establish an incremental fetch session.
LinkedHashMap<TopicPartition, FetchResponse.PartitionData<MemoryRecords>> partitions1 = new LinkedHashMap<>();
Expand Down Expand Up @@ -3503,7 +3503,7 @@ public void testOffsetValidationAwaitsNodeApiVersion() {

// Seek with a position and leader+epoch
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(
metadata.currentLeaderOrEmpty(tp0).leader, Optional.of(epochOne));
metadata.currentLeader(tp0).leader, Optional.of(epochOne));
subscriptions.seekUnvalidated(tp0, new SubscriptionState.FetchPosition(20L, Optional.of(epochOne), leaderAndEpoch));
assertFalse(client.isConnected(node.idString()));
assertTrue(subscriptions.awaitingValidation(tp0));
Expand Down Expand Up @@ -3552,7 +3552,7 @@ public void testOffsetValidationSkippedForOldBroker() {

// Seek with a position and leader+epoch
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(
metadata.currentLeaderOrEmpty(tp0).leader, Optional.of(epochOne));
metadata.currentLeader(tp0).leader, Optional.of(epochOne));
subscriptions.seekUnvalidated(tp0, new SubscriptionState.FetchPosition(0, Optional.of(epochOne), leaderAndEpoch));

// Update metadata to epoch=2, enter validation
Expand Down Expand Up @@ -3580,7 +3580,7 @@ public void testOffsetValidationHandlesSeekWithInflightOffsetForLeaderRequest()
Node node = metadata.fetch().nodes().get(0);
apiVersions.update(node.idString(), NodeApiVersions.create());

Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(metadata.currentLeaderOrEmpty(tp0).leader, Optional.of(epochOne));
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(metadata.currentLeader(tp0).leader, Optional.of(epochOne));
subscriptions.seekUnvalidated(tp0, new SubscriptionState.FetchPosition(0, Optional.of(epochOne), leaderAndEpoch));

fetcher.validateOffsetsIfNeeded();
Expand Down Expand Up @@ -3623,7 +3623,7 @@ public void testOffsetValidationFencing() {
apiVersions.update(node.idString(), NodeApiVersions.create());

// Seek with a position and leader+epoch
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(metadata.currentLeaderOrEmpty(tp0).leader, Optional.of(epochOne));
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(metadata.currentLeader(tp0).leader, Optional.of(epochOne));
subscriptions.seekValidated(tp0, new SubscriptionState.FetchPosition(0, Optional.of(epochOne), leaderAndEpoch));

// Update metadata to epoch=2, enter validation
Expand Down Expand Up @@ -3693,7 +3693,7 @@ public void testTruncationDetected() {
apiVersions.update(node.idString(), NodeApiVersions.create());

// Seek
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(metadata.currentLeaderOrEmpty(tp0).leader, Optional.of(1));
Metadata.LeaderAndEpoch leaderAndEpoch = new Metadata.LeaderAndEpoch(metadata.currentLeader(tp0).leader, Optional.of(1));
subscriptions.seekValidated(tp0, new SubscriptionState.FetchPosition(0, Optional.of(1), leaderAndEpoch));

// Check for truncation, this should cause tp0 to go into validation
Expand Down Expand Up @@ -3817,19 +3817,22 @@ public void testFetchCompletedBeforeHandlerAdded() {
client.prepareResponse(fullFetchResponse(tp0, buildRecords(1L, 1, 1), Errors.NONE, 100L, 0));
consumerClient.poll(time.timer(0));
fetchedRecords();
Node node = fetcher.selectReadReplica(tp0, subscriptions.position(tp0).currentLeader.leader, time.milliseconds());

Metadata.LeaderAndEpoch leaderAndEpoch = subscriptions.position(tp0).currentLeader;
assertTrue(leaderAndEpoch.leader.isPresent());
Node readReplica = fetcher.selectReadReplica(tp0, leaderAndEpoch.leader.get(), time.milliseconds());

AtomicBoolean wokenUp = new AtomicBoolean(false);
client.setWakeupHook(() -> {
if (!wokenUp.getAndSet(true)) {
consumerClient.disconnectAsync(node);
consumerClient.disconnectAsync(readReplica);
consumerClient.poll(time.timer(0));
}
});

assertEquals(1, fetcher.sendFetches());

consumerClient.disconnectAsync(node);
consumerClient.disconnectAsync(readReplica);
consumerClient.poll(time.timer(0));

assertEquals(1, fetcher.sendFetches());
Expand Down
Loading