From f146139a88d84ec2a95e45db8eb7ab94756e5911 Mon Sep 17 00:00:00 2001 From: Niket Goel Date: Tue, 24 May 2022 14:39:33 -0700 Subject: [PATCH 1/3] KAFKA-13888: Addition of Information in DescribeQuorumResponse about Voter Lag This commit adds an Admin API handler for DescribeQuorum Request and also adds in two new fields LastFetchTimestamp and LastCaughtUpTimestamp to the DescribeQuorumResponse as described by KIP-836. This commit does not implement the newly added fields. Those will be added in a subsequent commit. --- .../org/apache/kafka/clients/admin/Admin.java | 17 ++++ .../clients/admin/DescribeQuorumOptions.java | 25 ++++++ .../clients/admin/DescribeQuorumResult.java | 40 +++++++++ .../kafka/clients/admin/KafkaAdminClient.java | 67 ++++++++++++++ .../requests/DescribeQuorumResponse.java | 48 ++++++++++ .../apache/kafka/common/utils/QuorumInfo.java | 90 +++++++++++++++++++ .../common/message/DescribeQuorumRequest.json | 2 +- .../message/DescribeQuorumResponse.json | 8 +- .../kafka/clients/admin/MockAdminClient.java | 5 ++ .../kafka/admin/DescribeQuorumTest.scala | 64 +++++++++++++ .../org/apache/kafka/raft/RaftConfig.java | 15 ++++ 11 files changed, 378 insertions(+), 3 deletions(-) create mode 100644 clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java create mode 100644 clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java create mode 100644 clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java create mode 100644 core/src/test/scala/integration/kafka/admin/DescribeQuorumTest.scala diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java index 0c795bc5206dc..26ffd355452a0 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java @@ -1382,6 +1382,23 @@ default DescribeFeaturesResult describeFeatures() { return describeFeatures(new DescribeFeaturesOptions()); } + /** + * Describe the state of the raft quorum + *

+ * The following exceptions can be anticipated when calling {@code get()} on the futures obtained from + * the returned {@code DescribeQuorumResult}: + *

+ * + * @param options The options to use when describing the quorum. + * @return The DescribeQuorumResult. + */ + DescribeQuorumResult describeQuorum(DescribeQuorumOptions options); + /** * Describes finalized as well as supported features. The request is issued to any random * broker. diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java new file mode 100644 index 0000000000000..1fae819ee4aa7 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java @@ -0,0 +1,25 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.clients.admin; + +/** + * Options for {@link ConfluentAdmin#describeQuorum(DescribeQuorumOptions)}. + * + */ +public class DescribeQuorumOptions extends AbstractOptions { + +} diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java new file mode 100644 index 0000000000000..e4b94f8cde64c --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java @@ -0,0 +1,40 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.clients.admin; + +import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.common.QuorumInfo; + +/** + * The result of {@link Admin#describeQuorum(DescribeQuorumOptions)} + * + */ +public class DescribeQuorumResult { + + private final KafkaFuture quorumInfo; + + public DescribeQuorumResult(KafkaFuture quorumInfo) { + this.quorumInfo = quorumInfo; + } + + /** + * Returns a future QuorumInfo + */ + public KafkaFuture quorumInfo() { + return quorumInfo; + } +} diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index cc913bfda0e75..5cde2ee3ed7df 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -63,6 +63,8 @@ import org.apache.kafka.common.MetricName; import org.apache.kafka.common.Node; import org.apache.kafka.common.PartitionInfo; +import org.apache.kafka.common.QuorumInfo; +import org.apache.kafka.common.QuorumInfo.ReplicaState; import org.apache.kafka.common.TopicCollection; import org.apache.kafka.common.TopicCollection.TopicIdCollection; import org.apache.kafka.common.TopicCollection.TopicNameCollection; @@ -135,6 +137,7 @@ import org.apache.kafka.common.message.DescribeLogDirsRequestData; import org.apache.kafka.common.message.DescribeLogDirsRequestData.DescribableLogDirTopic; import org.apache.kafka.common.message.DescribeLogDirsResponseData; +import org.apache.kafka.common.message.DescribeQuorumRequestData; import org.apache.kafka.common.message.DescribeUserScramCredentialsRequestData; import org.apache.kafka.common.message.DescribeUserScramCredentialsRequestData.UserName; import org.apache.kafka.common.message.DescribeUserScramCredentialsResponseData; @@ -208,6 +211,9 @@ import org.apache.kafka.common.requests.DescribeLogDirsResponse; import org.apache.kafka.common.requests.DescribeUserScramCredentialsRequest; import org.apache.kafka.common.requests.DescribeUserScramCredentialsResponse; +import org.apache.kafka.common.requests.DescribeQuorumRequest; +import org.apache.kafka.common.requests.DescribeQuorumRequest.Builder; +import org.apache.kafka.common.requests.DescribeQuorumResponse; import org.apache.kafka.common.requests.ElectLeadersRequest; import org.apache.kafka.common.requests.ElectLeadersResponse; import org.apache.kafka.common.requests.ExpireDelegationTokenRequest; @@ -268,6 +274,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; +import static java.util.Collections.singletonList; import static org.apache.kafka.common.message.AlterPartitionReassignmentsRequestData.ReassignablePartition; import static org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignablePartitionResponse; import static org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignableTopicResponse; @@ -318,6 +325,11 @@ public class KafkaAdminClient extends AdminClient { private final Logger log; private final LogContext logContext; + /** + * The name of the internal raft metadata topic + */ + private static final String METADATA_TOPIC_NAME = "__cluster_metadata"; + /** * The default timeout to use for an operation. */ @@ -4181,6 +4193,61 @@ private static byte[] getSaltedPasword(ScramMechanism publicScramMechanism, byte .hi(password, salt, iterations); } + @Override + public DescribeQuorumResult describeQuorum(DescribeQuorumOptions options) { + NodeProvider provider = new LeastLoadedNodeProvider(); + + final KafkaFutureImpl future = new KafkaFutureImpl<>(); + final long now = time.milliseconds(); + final Call call = new Call( + "describeQuorum", calcDeadlineMs(now, options.timeoutMs()), provider) { + + private QuorumInfo createQuorumResult(final DescribeQuorumResponse response) { + Integer partition = 0; + String topicName = response.getTopicNameByIndex(partition); + Integer leaderId = response.getPartitionLeaderId(topicName, partition); + List voters = new ArrayList<>(); + for (Map.Entry entry: response.getVoterOffsets(topicName, partition).entrySet()) { + voters.add(new ReplicaState(entry.getKey(), entry.getValue())); + } + List observers = new ArrayList<>(); + for (Map.Entry entry: response.getObserverOffsets(topicName, partition).entrySet()) { + observers.add(new ReplicaState(entry.getKey(), entry.getValue())); + } + QuorumInfo info = new QuorumInfo(topicName, leaderId, voters, observers); + return info; + } + + @Override + DescribeQuorumRequest.Builder createRequest(int timeoutMs) { + DescribeQuorumRequestData data = new DescribeQuorumRequestData() + .setTopics(singletonList(new DescribeQuorumRequestData.TopicData() + .setPartitions(singletonList(new DescribeQuorumRequestData.PartitionData() + .setPartitionIndex(0))) + .setTopicName(METADATA_TOPIC_NAME))); + return new Builder(data); + } + + @Override + void handleResponse(AbstractResponse response) { + final DescribeQuorumResponse quorumResponse = (DescribeQuorumResponse) response; + if (quorumResponse.data().errorCode() == Errors.NONE.code()) { + future.complete(createQuorumResult(quorumResponse)); + } else { + future.completeExceptionally(Errors.forCode(quorumResponse.data().errorCode()).exception()); + } + } + + @Override + void handleFailure(Throwable throwable) { + completeAllExceptionally(Collections.singletonList(future), throwable); + } + }; + + runnable.call(call, now); + return new DescribeQuorumResult(future); + } + @Override public DescribeFeaturesResult describeFeatures(final DescribeFeaturesOptions options) { final KafkaFutureImpl future = new KafkaFutureImpl<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java index cbf945b70409a..5cdba86194f5d 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.message.DescribeQuorumResponseData; import org.apache.kafka.common.message.DescribeQuorumResponseData.ReplicaState; +import org.apache.kafka.common.message.DescribeQuorumResponseData.TopicData; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.ByteBufferAccessor; import org.apache.kafka.common.protocol.Errors; @@ -93,4 +94,51 @@ public static DescribeQuorumResponseData singletonResponse(TopicPartition topicP public static DescribeQuorumResponse parse(ByteBuffer buffer, short version) { return new DescribeQuorumResponse(new DescribeQuorumResponseData(new ByteBufferAccessor(buffer), version)); } + + public String getTopicNameByIndex(Integer index) { + return data.topics().get(index).topicName(); + } + + public Integer getPartitionLeaderId(String topicName, Integer partition) { + Integer leaderId = -1; + TopicData topic = data.topics().stream() + .filter(t -> t.topicName().equals(topicName)) + .findFirst() + .orElse(null); + if (topic != null) { + leaderId = Integer.valueOf(topic.partitions().get(partition).leaderId()); + } + return leaderId; + } + + public Map getVoterOffsets(String topicName, Integer partition) { + Map voterOffsets = new HashMap<>(); + TopicData topic = data.topics().stream() + .filter(t -> t.topicName().equals(topicName)) + .findFirst() + .orElse(null); + if(topic != null) { + topic.partitions().get(partition).currentVoters().forEach( + v -> { + voterOffsets.put(v.replicaId(), v.logEndOffset()); + } + ); + } + return voterOffsets; + } + + public Map getObserverOffsets(String topicName, Integer partition) { + Map observerOffsets = new HashMap<>(); + TopicData topic = data.topics().stream() + .filter(t -> t.topicName().equals(topicName)) + .findFirst() + .orElse(null); + topic.partitions().get(partition).observers().forEach( + o -> { + observerOffsets.put(o.replicaId(), o.logEndOffset()); + } + ); + return observerOffsets; + } + } diff --git a/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java b/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java new file mode 100644 index 0000000000000..a69e75a13ac8a --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java @@ -0,0 +1,90 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.common; + +import java.util.List; + +/** + * This is used to describe per-partition state in the DescribeQuorumResponse. + */ +public class QuorumInfo { + private final String topic; + private final Integer leaderId; + private final List voters; + private final List observers; + + public QuorumInfo(String topic, Integer leaderId, List voters, List observers) { + this.topic = topic; + this.leaderId = leaderId; + this.voters = voters; + this.observers = observers; + } + + public String topic() { + return topic; + } + + public Integer leaderId() { + return leaderId; + } + + public List voters() { + return voters; + } + + public List observers() { + return observers; + } + + public static class ReplicaState { + private final int replicaId; + private final long logEndOffset; + private final long lastFetchTimeMs; + private final long lastCaughtUpTimeMs; + + public ReplicaState(int replicaId, long logEndOffset) { + this.replicaId = replicaId; + this.logEndOffset = logEndOffset; + this.lastFetchTimeMs = -1; + this.lastCaughtUpTimeMs = -1; + } + + public ReplicaState(int replicaId, long logEndOffset, + long lastFetchTimeMs, long lastCaughtUpTimeMs) { + this.replicaId = replicaId; + this.logEndOffset = logEndOffset; + this.lastFetchTimeMs = lastFetchTimeMs; + this.lastCaughtUpTimeMs = lastCaughtUpTimeMs; + } + + public int replicaId() { + return replicaId; + } + + public long logEndOffset() { + return logEndOffset; + } + + public long lastFetchTimeMs() { + return lastFetchTimeMs; + } + + public long lastCaughtUpTimeMs() { + return lastCaughtUpTimeMs; + } + } +} diff --git a/clients/src/main/resources/common/message/DescribeQuorumRequest.json b/clients/src/main/resources/common/message/DescribeQuorumRequest.json index cd4a7f1db5470..6095763d594e4 100644 --- a/clients/src/main/resources/common/message/DescribeQuorumRequest.json +++ b/clients/src/main/resources/common/message/DescribeQuorumRequest.json @@ -18,7 +18,7 @@ "type": "request", "listeners": ["broker", "controller"], "name": "DescribeQuorumRequest", - "validVersions": "0", + "validVersions": "0-1", "flexibleVersions": "0+", "fields": [ { "name": "Topics", "type": "[]TopicData", diff --git a/clients/src/main/resources/common/message/DescribeQuorumResponse.json b/clients/src/main/resources/common/message/DescribeQuorumResponse.json index 444fee355a8ba..3fa993f410cf7 100644 --- a/clients/src/main/resources/common/message/DescribeQuorumResponse.json +++ b/clients/src/main/resources/common/message/DescribeQuorumResponse.json @@ -17,7 +17,7 @@ "apiKey": 55, "type": "response", "name": "DescribeQuorumResponse", - "validVersions": "0", + "validVersions": "0-1", "flexibleVersions": "0+", "fields": [ { "name": "ErrorCode", "type": "int16", "versions": "0+", @@ -44,7 +44,11 @@ { "name": "ReplicaState", "versions": "0+", "fields": [ { "name": "ReplicaId", "type": "int32", "versions": "0+", "entityType": "brokerId" }, { "name": "LogEndOffset", "type": "int64", "versions": "0+", - "about": "The last known log end offset of the follower or -1 if it is unknown"} + "about": "The last known log end offset of the follower or -1 if it is unknown"}, + { "name": "LastFetchTimestamp", "type": "int64", "versions": "1+", "ignorable": true, "default": -1, + "about": "The last known leader wall clock time time when a follower fetched from the leader. This is reported as -1 both for the current leader or if it is unknown for a voter"}, + { "name": "LastCaughtUpTimestamp", "type": "int64", "versions": "1+", "ignorable": true, "default": -1, + "about": "The leader wall clock append time of the offset for which the follower made the most recent fetch request. This is reported as the current time for the leader and -1 if unknown for a voter"} ]} ] } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java b/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java index 15cdc5ccc4116..d257aa8635c24 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java @@ -962,6 +962,11 @@ public AlterUserScramCredentialsResult alterUserScramCredentials(List cluster.brokers.get(i).brokerState == BrokerState.RUNNING, + "Broker Never started up") + } + val props = cluster.clientProperties() + props.put(AdminClientConfig.CLIENT_ID_CONFIG, this.getClass.getName) + val admin = Admin.create(props) + try { + val quorumState = admin.describeQuorum(new DescribeQuorumOptions) + val quorumInfo = quorumState.quorumInfo().get() + + assertEquals(KafkaRaftServer.MetadataTopic, quorumInfo.topic()) + assertEquals(3, quorumInfo.voters.size()) + } finally { + admin.close() + } + } finally { + cluster.close + } + } +} diff --git a/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java b/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java index 3ce72a591fe1e..8f1aeb7a880a7 100644 --- a/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java +++ b/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java @@ -256,6 +256,21 @@ public static List voterConnectionsToNodes(Map nodes) { + StringBuilder voterConnections = new StringBuilder(""); + nodes.forEach(node -> { + if (!voterConnections.toString().isEmpty()) { + voterConnections.append(","); + } + voterConnections.append(node.idString()); + voterConnections.append("@"); + voterConnections.append(node.host()); + voterConnections.append(":"); + voterConnections.append(node.port()); + }); + return voterConnections.toString(); + } + public static class ControllerQuorumVotersValidator implements ConfigDef.Validator { @Override public void ensureValid(String name, Object value) { From c3dfe9d5b0a1905023f0599091101ea1906c5200 Mon Sep 17 00:00:00 2001 From: Niket Goel Date: Tue, 24 May 2022 19:59:08 -0700 Subject: [PATCH 2/3] Addressing review comments Changes: Added some unit tests to KafkaAdminClientTest Merged multiple DescribeQuorumTest files Some other refactoring --- .../org/apache/kafka/clients/admin/Admin.java | 46 +++-- ...ava => DescribeMetadataQuorumOptions.java} | 4 +- ...java => DescribeMetadataQuorumResult.java} | 24 +-- .../kafka/clients/admin/KafkaAdminClient.java | 117 ++++++------ .../apache/kafka/common/internals/Topic.java | 1 + .../requests/DescribeQuorumResponse.java | 4 +- .../apache/kafka/common/utils/QuorumInfo.java | 167 ++++++++++++------ .../message/DescribeQuorumResponse.json | 1 + .../clients/admin/KafkaAdminClientTest.java | 51 ++++++ .../kafka/clients/admin/MockAdminClient.java | 2 +- .../scala/kafka/server/KafkaRaftServer.scala | 3 +- .../kafka/admin/DescribeQuorumTest.scala | 64 ------- ...la => DescribeQuorumIntegrationTest.scala} | 43 ++++- .../org/apache/kafka/raft/RaftConfig.java | 15 -- 14 files changed, 310 insertions(+), 232 deletions(-) rename clients/src/main/java/org/apache/kafka/clients/admin/{DescribeQuorumOptions.java => DescribeMetadataQuorumOptions.java} (82%) rename clients/src/main/java/org/apache/kafka/clients/admin/{DescribeQuorumResult.java => DescribeMetadataQuorumResult.java} (68%) delete mode 100644 core/src/test/scala/integration/kafka/admin/DescribeQuorumTest.scala rename core/src/test/scala/unit/kafka/server/{DescribeQuorumRequestTest.scala => DescribeQuorumIntegrationTest.scala} (71%) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java index 26ffd355452a0..95eebedd1dcec 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java @@ -1382,23 +1382,6 @@ default DescribeFeaturesResult describeFeatures() { return describeFeatures(new DescribeFeaturesOptions()); } - /** - * Describe the state of the raft quorum - *

- * The following exceptions can be anticipated when calling {@code get()} on the futures obtained from - * the returned {@code DescribeQuorumResult}: - *

    - *
  • {@link org.apache.kafka.common.errors.ClusterAuthorizationException} - * If the authenticated user didn't have {@code DESCRIBE} access to the cluster.
  • - *
  • {@link org.apache.kafka.common.errors.TimeoutException} - * If the request timed out before the controller could list the cluster links.
  • - *
- * - * @param options The options to use when describing the quorum. - * @return The DescribeQuorumResult. - */ - DescribeQuorumResult describeQuorum(DescribeQuorumOptions options); - /** * Describes finalized as well as supported features. The request is issued to any random * broker. @@ -1463,6 +1446,35 @@ default DescribeFeaturesResult describeFeatures() { */ UpdateFeaturesResult updateFeatures(Map featureUpdates, UpdateFeaturesOptions options); + /** + * Describes the state of the metadata quorum. + *

+ * This is a convenience method for {@link #describeMetadataQuorum(DescribeMetadataQuorumOptions)} with default options. + * See the overload for more details. + * + * @return the {@link DescribeMetadataQuorumResult} containing the result + */ + default DescribeMetadataQuorumResult describeMetadataQuorum() { + return describeMetadataQuorum(new DescribeMetadataQuorumOptions()); + } + + /** + * Describes the state of the metadata quorum. + *

+ * The following exceptions can be anticipated when calling {@code get()} on the futures obtained from + * the returned {@code DescribeMetadataQuorumResult}: + *

    + *
  • {@link org.apache.kafka.common.errors.ClusterAuthorizationException} + * If the authenticated user didn't have {@code DESCRIBE} access to the cluster.
  • + *
  • {@link org.apache.kafka.common.errors.TimeoutException} + * If the request timed out before the controller could list the cluster links.
  • + *
+ * + * @param options The {@link DescribeMetadataQuorumOptions} to use when describing the quorum. + * @return the {@link DescribeMetadataQuorumResult} containing the result + */ + DescribeMetadataQuorumResult describeMetadataQuorum(DescribeMetadataQuorumOptions options); + /** * Unregister a broker. *

diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeMetadataQuorumOptions.java similarity index 82% rename from clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java rename to clients/src/main/java/org/apache/kafka/clients/admin/DescribeMetadataQuorumOptions.java index 1fae819ee4aa7..772f4eaa0aaa8 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumOptions.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeMetadataQuorumOptions.java @@ -17,9 +17,9 @@ package org.apache.kafka.clients.admin; /** - * Options for {@link ConfluentAdmin#describeQuorum(DescribeQuorumOptions)}. + * Options for {@link ConfluentAdmin#describeQuorum(DescribeMetadataQuorumOptions)}. * */ -public class DescribeQuorumOptions extends AbstractOptions { +public class DescribeMetadataQuorumOptions extends AbstractOptions { } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeMetadataQuorumResult.java similarity index 68% rename from clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java rename to clients/src/main/java/org/apache/kafka/clients/admin/DescribeMetadataQuorumResult.java index e4b94f8cde64c..a29030d930790 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/DescribeQuorumResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/DescribeMetadataQuorumResult.java @@ -20,21 +20,21 @@ import org.apache.kafka.common.QuorumInfo; /** - * The result of {@link Admin#describeQuorum(DescribeQuorumOptions)} + * The result of {@link Admin#describeMetadataQuorum(DescribeMetadataQuorumOptions)} * */ -public class DescribeQuorumResult { +public class DescribeMetadataQuorumResult { - private final KafkaFuture quorumInfo; + private final KafkaFuture quorumInfo; - public DescribeQuorumResult(KafkaFuture quorumInfo) { - this.quorumInfo = quorumInfo; - } + public DescribeMetadataQuorumResult(KafkaFuture quorumInfo) { + this.quorumInfo = quorumInfo; + } - /** - * Returns a future QuorumInfo - */ - public KafkaFuture quorumInfo() { - return quorumInfo; - } + /** + * Returns a future QuorumInfo + */ + public KafkaFuture quorumInfo() { + return quorumInfo; + } } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 5cde2ee3ed7df..0d19a286081c1 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -263,6 +263,7 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.OptionalLong; import java.util.Set; import java.util.TreeMap; import java.util.concurrent.TimeUnit; @@ -275,6 +276,7 @@ import java.util.stream.Stream; import static java.util.Collections.singletonList; +import static org.apache.kafka.common.internals.Topic.METADATA_TOPIC_NAME; import static org.apache.kafka.common.message.AlterPartitionReassignmentsRequestData.ReassignablePartition; import static org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignablePartitionResponse; import static org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignableTopicResponse; @@ -325,11 +327,6 @@ public class KafkaAdminClient extends AdminClient { private final Logger log; private final LogContext logContext; - /** - * The name of the internal raft metadata topic - */ - private static final String METADATA_TOPIC_NAME = "__cluster_metadata"; - /** * The default timeout to use for an operation. */ @@ -4193,61 +4190,6 @@ private static byte[] getSaltedPasword(ScramMechanism publicScramMechanism, byte .hi(password, salt, iterations); } - @Override - public DescribeQuorumResult describeQuorum(DescribeQuorumOptions options) { - NodeProvider provider = new LeastLoadedNodeProvider(); - - final KafkaFutureImpl future = new KafkaFutureImpl<>(); - final long now = time.milliseconds(); - final Call call = new Call( - "describeQuorum", calcDeadlineMs(now, options.timeoutMs()), provider) { - - private QuorumInfo createQuorumResult(final DescribeQuorumResponse response) { - Integer partition = 0; - String topicName = response.getTopicNameByIndex(partition); - Integer leaderId = response.getPartitionLeaderId(topicName, partition); - List voters = new ArrayList<>(); - for (Map.Entry entry: response.getVoterOffsets(topicName, partition).entrySet()) { - voters.add(new ReplicaState(entry.getKey(), entry.getValue())); - } - List observers = new ArrayList<>(); - for (Map.Entry entry: response.getObserverOffsets(topicName, partition).entrySet()) { - observers.add(new ReplicaState(entry.getKey(), entry.getValue())); - } - QuorumInfo info = new QuorumInfo(topicName, leaderId, voters, observers); - return info; - } - - @Override - DescribeQuorumRequest.Builder createRequest(int timeoutMs) { - DescribeQuorumRequestData data = new DescribeQuorumRequestData() - .setTopics(singletonList(new DescribeQuorumRequestData.TopicData() - .setPartitions(singletonList(new DescribeQuorumRequestData.PartitionData() - .setPartitionIndex(0))) - .setTopicName(METADATA_TOPIC_NAME))); - return new Builder(data); - } - - @Override - void handleResponse(AbstractResponse response) { - final DescribeQuorumResponse quorumResponse = (DescribeQuorumResponse) response; - if (quorumResponse.data().errorCode() == Errors.NONE.code()) { - future.complete(createQuorumResult(quorumResponse)); - } else { - future.completeExceptionally(Errors.forCode(quorumResponse.data().errorCode()).exception()); - } - } - - @Override - void handleFailure(Throwable throwable) { - completeAllExceptionally(Collections.singletonList(future), throwable); - } - }; - - runnable.call(call, now); - return new DescribeQuorumResult(future); - } - @Override public DescribeFeaturesResult describeFeatures(final DescribeFeaturesOptions options) { final KafkaFutureImpl future = new KafkaFutureImpl<>(); @@ -4388,6 +4330,61 @@ void handleFailure(Throwable throwable) { return new UpdateFeaturesResult(new HashMap<>(updateFutures)); } + @Override + public DescribeMetadataQuorumResult describeMetadataQuorum(DescribeMetadataQuorumOptions options) { + NodeProvider provider = new LeastLoadedNodeProvider(); + + final KafkaFutureImpl future = new KafkaFutureImpl<>(); + final long now = time.milliseconds(); + final Call call = new Call( + "describeQuorum", calcDeadlineMs(now, options.timeoutMs()), provider) { + + private QuorumInfo createQuorumResult(final DescribeQuorumResponse response) { + Integer partition = 0; + String topicName = response.getTopicNameByIndex(partition); + Integer leaderId = response.getPartitionLeaderId(topicName, partition); + List voters = new ArrayList<>(); + for (Map.Entry entry: response.getVoterOffsets(topicName, partition).entrySet()) { + voters.add(new ReplicaState(entry.getKey(), entry.getValue(), OptionalLong.empty(), OptionalLong.empty())); + } + List observers = new ArrayList<>(); + for (Map.Entry entry: response.getObserverOffsets(topicName, partition).entrySet()) { + observers.add(new ReplicaState(entry.getKey(), entry.getValue(), OptionalLong.empty(), OptionalLong.empty())); + } + QuorumInfo info = new QuorumInfo(topicName, leaderId, voters, observers); + return info; + } + + @Override + DescribeQuorumRequest.Builder createRequest(int timeoutMs) { + DescribeQuorumRequestData data = new DescribeQuorumRequestData() + .setTopics(singletonList(new DescribeQuorumRequestData.TopicData() + .setPartitions(singletonList(new DescribeQuorumRequestData.PartitionData() + .setPartitionIndex(0))) + .setTopicName(METADATA_TOPIC_NAME))); + return new Builder(data); + } + + @Override + void handleResponse(AbstractResponse response) { + final DescribeQuorumResponse quorumResponse = (DescribeQuorumResponse) response; + if (quorumResponse.data().errorCode() == Errors.NONE.code()) { + future.complete(createQuorumResult(quorumResponse)); + } else { + future.completeExceptionally(Errors.forCode(quorumResponse.data().errorCode()).exception()); + } + } + + @Override + void handleFailure(Throwable throwable) { + completeAllExceptionally(Collections.singletonList(future), throwable); + } + }; + + runnable.call(call, now); + return new DescribeMetadataQuorumResult(future); + } + @Override public UnregisterBrokerResult unregisterBroker(int brokerId, UnregisterBrokerOptions options) { final KafkaFutureImpl future = new KafkaFutureImpl<>(); diff --git a/clients/src/main/java/org/apache/kafka/common/internals/Topic.java b/clients/src/main/java/org/apache/kafka/common/internals/Topic.java index 3c93ef87b5c99..b23c9241fba70 100644 --- a/clients/src/main/java/org/apache/kafka/common/internals/Topic.java +++ b/clients/src/main/java/org/apache/kafka/common/internals/Topic.java @@ -27,6 +27,7 @@ public class Topic { public static final String GROUP_METADATA_TOPIC_NAME = "__consumer_offsets"; public static final String TRANSACTION_STATE_TOPIC_NAME = "__transaction_state"; + public static final String METADATA_TOPIC_NAME = "__cluster_metadata"; public static final String LEGAL_CHARS = "[a-zA-Z0-9._-]"; private static final Set INTERNAL_TOPICS = Collections.unmodifiableSet( diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java index 5cdba86194f5d..a97dafc009f11 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java @@ -106,7 +106,7 @@ public Integer getPartitionLeaderId(String topicName, Integer partition) { .findFirst() .orElse(null); if (topic != null) { - leaderId = Integer.valueOf(topic.partitions().get(partition).leaderId()); + leaderId = Integer.valueOf(topic.partitions().get(partition).leaderId()); } return leaderId; } @@ -117,7 +117,7 @@ public Map getVoterOffsets(String topicName, Integer partition) { .filter(t -> t.topicName().equals(topicName)) .findFirst() .orElse(null); - if(topic != null) { + if (topic != null) { topic.partitions().get(partition).currentVoters().forEach( v -> { voterOffsets.put(v.replicaId(), v.logEndOffset()); diff --git a/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java b/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java index a69e75a13ac8a..5ba8f7b05af2c 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/QuorumInfo.java @@ -17,74 +17,133 @@ package org.apache.kafka.common; import java.util.List; +import java.util.Objects; +import java.util.OptionalLong; /** * This is used to describe per-partition state in the DescribeQuorumResponse. */ public class QuorumInfo { - private final String topic; - private final Integer leaderId; - private final List voters; - private final List observers; - - public QuorumInfo(String topic, Integer leaderId, List voters, List observers) { - this.topic = topic; - this.leaderId = leaderId; - this.voters = voters; - this.observers = observers; - } - - public String topic() { - return topic; - } - - public Integer leaderId() { - return leaderId; - } - - public List voters() { - return voters; - } - - public List observers() { - return observers; - } - - public static class ReplicaState { - private final int replicaId; - private final long logEndOffset; - private final long lastFetchTimeMs; - private final long lastCaughtUpTimeMs; - - public ReplicaState(int replicaId, long logEndOffset) { - this.replicaId = replicaId; - this.logEndOffset = logEndOffset; - this.lastFetchTimeMs = -1; - this.lastCaughtUpTimeMs = -1; + private final String topic; + private final Integer leaderId; + private final List voters; + private final List observers; + + public QuorumInfo(String topic, Integer leaderId, List voters, List observers) { + this.topic = topic; + this.leaderId = leaderId; + this.voters = voters; + this.observers = observers; + } + + public String topic() { + return topic; + } + + public Integer leaderId() { + return leaderId; + } + + public List voters() { + return voters; } - public ReplicaState(int replicaId, long logEndOffset, - long lastFetchTimeMs, long lastCaughtUpTimeMs) { - this.replicaId = replicaId; - this.logEndOffset = logEndOffset; - this.lastFetchTimeMs = lastFetchTimeMs; - this.lastCaughtUpTimeMs = lastCaughtUpTimeMs; + public List observers() { + return observers; } - public int replicaId() { - return replicaId; + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + QuorumInfo that = (QuorumInfo) o; + return topic.equals(that.topic) + && leaderId.equals(that.leaderId) + && voters.equals(that.voters) + && observers.equals(that.observers); } - public long logEndOffset() { - return logEndOffset; + @Override + public int hashCode() { + return Objects.hash(topic, leaderId, voters, observers); } - public long lastFetchTimeMs() { - return lastFetchTimeMs; + @Override + public String toString() { + return "QuorumInfo{" + + "topic='" + topic + '\'' + + ", leaderId=" + leaderId + + ", voters=" + voters.toString() + + ", observers=" + observers.toString() + + '}'; } - public long lastCaughtUpTimeMs() { - return lastCaughtUpTimeMs; + public static class ReplicaState { + private final int replicaId; + private final long logEndOffset; + private final OptionalLong lastFetchTimeMs; + private final OptionalLong lastCaughtUpTimeMs; + + public ReplicaState() { + this(0, 0, OptionalLong.empty(), OptionalLong.empty()); + } + + public ReplicaState(int replicaId, long logEndOffset, + OptionalLong lastFetchTimeMs, OptionalLong lastCaughtUpTimeMs) { + this.replicaId = replicaId; + this.logEndOffset = logEndOffset; + this.lastFetchTimeMs = lastFetchTimeMs; + this.lastCaughtUpTimeMs = lastCaughtUpTimeMs; + } + + public int replicaId() { + return replicaId; + } + + public long logEndOffset() { + return logEndOffset; + } + + /** + * Return the lastFetchTime in milliseconds for this replica. + * @return The value of the lastFetchTime if known, empty otherwise + */ + public OptionalLong lastFetchTimeMs() { + return lastFetchTimeMs; + } + + /** + * Return the lastCaughtUpTime in milliseconds for this replica. + * @return The value of the lastCaughtUpTime if known, empty otherwise + */ + public OptionalLong lastCaughtUpTimeMs() { + return lastCaughtUpTimeMs; + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + ReplicaState that = (ReplicaState) o; + return replicaId == that.replicaId + && logEndOffset == that.logEndOffset + && lastFetchTimeMs.equals(that.lastFetchTimeMs) + && lastCaughtUpTimeMs.equals(that.lastCaughtUpTimeMs); + } + + @Override + public int hashCode() { + return Objects.hash(replicaId, logEndOffset, lastFetchTimeMs, lastCaughtUpTimeMs); + } + + @Override + public String toString() { + return "ReplicaState{" + + "replicaId=" + replicaId + + ", logEndOffset=" + logEndOffset + + ", lastFetchTimeMs=" + lastFetchTimeMs + + ", lastCaughtUpTimeMs=" + lastCaughtUpTimeMs + + '}'; + } } - } } diff --git a/clients/src/main/resources/common/message/DescribeQuorumResponse.json b/clients/src/main/resources/common/message/DescribeQuorumResponse.json index 3fa993f410cf7..e94b3327d695c 100644 --- a/clients/src/main/resources/common/message/DescribeQuorumResponse.json +++ b/clients/src/main/resources/common/message/DescribeQuorumResponse.json @@ -20,6 +20,7 @@ "validVersions": "0-1", "flexibleVersions": "0+", "fields": [ + // Version 1 adds LastFetchTimeStamp and LastCaughtUpTimestamp in ReplicaState { "name": "ErrorCode", "type": "int16", "versions": "0+", "about": "The top level error code."}, { "name": "Topics", "type": "[]TopicData", diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index eb4681856e273..4b5238c6738df 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -31,6 +31,7 @@ import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.Node; +import org.apache.kafka.common.QuorumInfo; import org.apache.kafka.common.TopicCollection; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; @@ -69,6 +70,7 @@ import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.feature.Features; +import org.apache.kafka.common.internals.Topic; import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData; import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData; import org.apache.kafka.common.message.AlterReplicaLogDirsResponseData.AlterReplicaLogDirPartitionResult; @@ -100,6 +102,7 @@ import org.apache.kafka.common.message.DescribeLogDirsResponseData; import org.apache.kafka.common.message.DescribeLogDirsResponseData.DescribeLogDirsTopic; import org.apache.kafka.common.message.DescribeProducersResponseData; +import org.apache.kafka.common.message.DescribeQuorumResponseData; import org.apache.kafka.common.message.DescribeTransactionsResponseData; import org.apache.kafka.common.message.DescribeUserScramCredentialsResponseData; import org.apache.kafka.common.message.DescribeUserScramCredentialsResponseData.CredentialInfo; @@ -161,6 +164,8 @@ import org.apache.kafka.common.requests.DescribeLogDirsResponse; import org.apache.kafka.common.requests.DescribeProducersRequest; import org.apache.kafka.common.requests.DescribeProducersResponse; +import org.apache.kafka.common.requests.DescribeQuorumRequest; +import org.apache.kafka.common.requests.DescribeQuorumResponse; import org.apache.kafka.common.requests.DescribeTransactionsRequest; import org.apache.kafka.common.requests.DescribeTransactionsResponse; import org.apache.kafka.common.requests.DescribeUserScramCredentialsResponse; @@ -592,6 +597,26 @@ private static ApiVersionsResponse prepareApiVersionsResponseForDescribeFeatures .setErrorCode(error.code())); } + private static QuorumInfo defaultQuorumInfo() { + return new QuorumInfo(Topic.METADATA_TOPIC_NAME, 0, + singletonList(new QuorumInfo.ReplicaState()), + singletonList(new QuorumInfo.ReplicaState())); + } + + private static DescribeQuorumResponse prepareDescribeQuorumResponse(Errors error) { + if (error == Errors.NONE) { + return new DescribeQuorumResponse(DescribeQuorumResponse.singletonResponse( + new TopicPartition(Topic.METADATA_TOPIC_NAME, 0), + 0, 0, 0, + singletonList(new DescribeQuorumResponseData.ReplicaState()), + singletonList(new DescribeQuorumResponseData.ReplicaState())) + ); + } + return new DescribeQuorumResponse( + new DescribeQuorumResponseData() + .setErrorCode(error.code())); + } + /** * Test that the client properly times out when we don't receive any metadata. */ @@ -4846,6 +4871,32 @@ public void testDescribeFeaturesFailure() { } } + @Test + public void testDescribeMetadataQuorumSuccess() throws Exception { + try (final AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create((short) 55, (short) 0, (short) 1)); + env.kafkaClient().prepareResponse( + body -> body instanceof DescribeQuorumRequest, + prepareDescribeQuorumResponse(Errors.NONE)); + final KafkaFuture future = env.adminClient().describeMetadataQuorum().quorumInfo(); + final QuorumInfo quorumInfo = future.get(); + assertEquals(defaultQuorumInfo(), quorumInfo); + } + } + + @Test + public void testDescribeMetadataQuorumFailure() { + try (final AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create((short) 55, (short) 0, (short) 1)); + env.kafkaClient().prepareResponse( + body -> body instanceof DescribeQuorumRequest, + prepareDescribeQuorumResponse(Errors.INVALID_REQUEST)); + final KafkaFuture future = env.adminClient().describeMetadataQuorum().quorumInfo(); + final ExecutionException e = assertThrows(ExecutionException.class, future::get); + assertEquals(e.getCause().getClass(), Errors.INVALID_REQUEST.exception().getClass()); + } + } + @Test public void testListOffsetsMetadataRetriableErrors() throws Exception { diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java b/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java index d257aa8635c24..7bac8803fe021 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java @@ -963,7 +963,7 @@ public AlterUserScramCredentialsResult alterUserScramCredentials(List cluster.brokers.get(i).brokerState == BrokerState.RUNNING, - "Broker Never started up") - } - val props = cluster.clientProperties() - props.put(AdminClientConfig.CLIENT_ID_CONFIG, this.getClass.getName) - val admin = Admin.create(props) - try { - val quorumState = admin.describeQuorum(new DescribeQuorumOptions) - val quorumInfo = quorumState.quorumInfo().get() - - assertEquals(KafkaRaftServer.MetadataTopic, quorumInfo.topic()) - assertEquals(3, quorumInfo.voters.size()) - } finally { - admin.close() - } - } finally { - cluster.close - } - } -} diff --git a/core/src/test/scala/unit/kafka/server/DescribeQuorumRequestTest.scala b/core/src/test/scala/unit/kafka/server/DescribeQuorumIntegrationTest.scala similarity index 71% rename from core/src/test/scala/unit/kafka/server/DescribeQuorumRequestTest.scala rename to core/src/test/scala/unit/kafka/server/DescribeQuorumIntegrationTest.scala index 55b9fe92c3c3a..aecb6a6ba2d89 100644 --- a/core/src/test/scala/unit/kafka/server/DescribeQuorumRequestTest.scala +++ b/core/src/test/scala/unit/kafka/server/DescribeQuorumIntegrationTest.scala @@ -17,25 +17,30 @@ package kafka.server import java.io.IOException - import kafka.test.ClusterInstance import kafka.test.annotation.{ClusterTest, ClusterTestDefaults, Type} import kafka.test.junit.ClusterTestExtensions -import kafka.utils.NotNothing +import kafka.testkit.{KafkaClusterTestKit, TestKitNodes} +import kafka.utils.{NotNothing, TestUtils} +import org.apache.kafka.clients.admin.{Admin, AdminClientConfig, DescribeMetadataQuorumOptions} import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.requests.DescribeQuorumRequest.singletonRequest import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse, ApiVersionsRequest, ApiVersionsResponse, DescribeQuorumRequest, DescribeQuorumResponse} +import org.apache.kafka.metadata.BrokerState import org.junit.jupiter.api.Assertions._ -import org.junit.jupiter.api.Tag +import org.junit.jupiter.api.{Tag, Timeout} import org.junit.jupiter.api.extension.ExtendWith +import org.slf4j.LoggerFactory import scala.jdk.CollectionConverters._ import scala.reflect.ClassTag +@Timeout(120) @ExtendWith(value = Array(classOf[ClusterTestExtensions])) @ClusterTestDefaults(clusterType = Type.KRAFT) @Tag("integration") -class DescribeQuorumRequestTest(cluster: ClusterInstance) { +class DescribeQuorumIntegrationTest(cluster: ClusterInstance) { + val log = LoggerFactory.getLogger(classOf[DescribeQuorumIntegrationTest]) @ClusterTest(clusterType = Type.ZK) def testDescribeQuorumNotSupportedByZkBrokers(): Unit = { @@ -80,6 +85,36 @@ class DescribeQuorumRequestTest(cluster: ClusterInstance) { assertTrue(leaderState.logEndOffset > 0) } + @ClusterTest + def testDescribeQuorumRequestToBrokers() = { + val cluster = new KafkaClusterTestKit.Builder( + new TestKitNodes.Builder(). + setNumBrokerNodes(4). + setNumControllerNodes(3).build()).build() + try { + cluster.format + cluster.startup + for (i <- 0 to 3) { + TestUtils.waitUntilTrue(() => cluster.brokers.get(i).brokerState == BrokerState.RUNNING, + "Broker Never started up") + } + val props = cluster.clientProperties() + props.put(AdminClientConfig.CLIENT_ID_CONFIG, this.getClass.getName) + val admin = Admin.create(props) + try { + val quorumState = admin.describeMetadataQuorum(new DescribeMetadataQuorumOptions) + val quorumInfo = quorumState.quorumInfo().get() + + assertEquals(KafkaRaftServer.MetadataTopic, quorumInfo.topic()) + assertEquals(3, quorumInfo.voters.size()) + } finally { + admin.close() + } + } finally { + cluster.close + } + } + private def connectAndReceive[T <: AbstractResponse]( request: AbstractRequest )( diff --git a/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java b/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java index 8f1aeb7a880a7..3ce72a591fe1e 100644 --- a/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java +++ b/raft/src/main/java/org/apache/kafka/raft/RaftConfig.java @@ -256,21 +256,6 @@ public static List voterConnectionsToNodes(Map nodes) { - StringBuilder voterConnections = new StringBuilder(""); - nodes.forEach(node -> { - if (!voterConnections.toString().isEmpty()) { - voterConnections.append(","); - } - voterConnections.append(node.idString()); - voterConnections.append("@"); - voterConnections.append(node.host()); - voterConnections.append(":"); - voterConnections.append(node.port()); - }); - return voterConnections.toString(); - } - public static class ControllerQuorumVotersValidator implements ConfigDef.Validator { @Override public void ensureValid(String name, Object value) { From e8f603c7d2e1934c8bc88dfe09579151b4e1419a Mon Sep 17 00:00:00 2001 From: lqjack Date: Thu, 26 May 2022 00:28:20 +0800 Subject: [PATCH 3/3] KAFKA-13888 create file if not exist update last --- .../apache/kafka/raft/KafkaRaftClient.java | 15 +++-- .../org/apache/kafka/raft/LeaderState.java | 55 ++++++++++++++++--- .../apache/kafka/raft/LeaderStateTest.java | 30 +++++----- .../raft/internals/KafkaRaftMetricsTest.java | 2 +- 4 files changed, 74 insertions(+), 28 deletions(-) diff --git a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java index 3479e05d32494..4cecb30e675aa 100644 --- a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java +++ b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java @@ -1011,7 +1011,7 @@ private FetchResponseData tryCompleteFetchRequest( if (validOffsetAndEpoch.kind() == ValidOffsetAndEpoch.Kind.VALID) { LogFetchInfo info = log.read(fetchOffset, Isolation.UNCOMMITTED); - if (state.updateReplicaState(replicaId, currentTimeMs, info.startOffsetMetadata)) { + if (state.updateReplicaState(replicaId, currentTimeMs, info.startOffsetMetadata, fetchOffset)) { onUpdateLeaderHighWatermark(state, currentTimeMs); } @@ -1174,8 +1174,8 @@ private DescribeQuorumResponseData handleDescribeQuorumRequest( leaderState.localId(), leaderState.epoch(), leaderState.highWatermark().isPresent() ? leaderState.highWatermark().get().offset : -1, - convertToReplicaStates(leaderState.getVoterEndOffsets()), - convertToReplicaStates(leaderState.getObserverStates(currentTimeMs)) + convertToReplicaStates(leaderState, true, currentTimeMs), + convertToReplicaStates(leaderState, false, currentTimeMs) ); } @@ -1412,12 +1412,17 @@ private boolean handleFetchSnapshotResponse( state.resetFetchTimeout(currentTimeMs); return true; } - - List convertToReplicaStates(Map replicaEndOffsets) { + @SuppressWarnings("unchecked") + List convertToReplicaStates(LeaderState leaderState, boolean voterOrObserver, long currentTimeMs) { + Map replicaEndOffsets = leaderState.getVoterEndOffsets(); + if (!voterOrObserver) { + replicaEndOffsets = leaderState.getObserverStates(currentTimeMs); + } return replicaEndOffsets.entrySet().stream() .map(entry -> new ReplicaState() .setReplicaId(entry.getKey()) .setLogEndOffset(entry.getValue())) + .peek(rep -> rep.setLastFetchTimestamp(leaderState.getReplicaLastFetchTimestamp(rep.replicaId()))) .collect(Collectors.toList()); } diff --git a/raft/src/main/java/org/apache/kafka/raft/LeaderState.java b/raft/src/main/java/org/apache/kafka/raft/LeaderState.java index de08b7b1cc983..f132d2297d378 100644 --- a/raft/src/main/java/org/apache/kafka/raft/LeaderState.java +++ b/raft/src/main/java/org/apache/kafka/raft/LeaderState.java @@ -205,10 +205,10 @@ private boolean updateHighWatermark() { /** * Update the local replica state. * - * See {@link #updateReplicaState(int, long, LogOffsetMetadata)} + * See {@link #updateReplicaState(int, long, LogOffsetMetadata, long)} */ public boolean updateLocalState(long fetchTimestamp, LogOffsetMetadata logOffsetMetadata) { - return updateReplicaState(localId, fetchTimestamp, logOffsetMetadata); + return updateReplicaState(localId, fetchTimestamp, logOffsetMetadata, -1L); } /** @@ -217,11 +217,13 @@ public boolean updateLocalState(long fetchTimestamp, LogOffsetMetadata logOffset * @param replicaId replica id * @param fetchTimestamp fetch timestamp * @param logOffsetMetadata new log offset and metadata + * @param fetchOffset fetch offset * @return true if the high watermark is updated too */ public boolean updateReplicaState(int replicaId, long fetchTimestamp, - LogOffsetMetadata logOffsetMetadata) { + LogOffsetMetadata logOffsetMetadata, + long fetchOffset) { // Ignore fetches from negative replica id, as it indicates // the fetch is from non-replica. For example, a consumer. if (replicaId < 0) { @@ -230,7 +232,7 @@ public boolean updateReplicaState(int replicaId, ReplicaState state = getReplicaState(replicaId); state.updateFetchTimestamp(fetchTimestamp); - return updateEndOffset(state, logOffsetMetadata); + return updateEndOffset(state, logOffsetMetadata, fetchTimestamp, fetchOffset); } public List nonLeaderVotersByDescendingFetchOffset() { @@ -246,8 +248,10 @@ private List followersByDescendingFetchOffset() { .collect(Collectors.toList()); } - private boolean updateEndOffset(ReplicaState state, - LogOffsetMetadata endOffsetMetadata) { + private boolean updateEndOffset(ReplicaState state, // leader status + LogOffsetMetadata endOffsetMetadata, //follower stats + long fetchTimestamp, + long fetchOffset) { state.endOffset.ifPresent(currentEndOffset -> { if (currentEndOffset.offset > endOffsetMetadata.offset) { if (state.nodeId == localId) { @@ -257,7 +261,13 @@ private boolean updateEndOffset(ReplicaState state, log.warn("Detected non-monotonic update of fetch offset from nodeId {}: {} -> {}", state.nodeId, currentEndOffset.offset, endOffsetMetadata.offset); } - } + } else if (endOffsetMetadata.offset >= currentEndOffset.offset) { + state.lastCaughtUpTimestamp = OptionalLong.of(Math.max(state.lastCaughtUpTimestamp.orElse(-1), fetchTimestamp)); + } else if (fetchOffset >= currentEndOffset.offset) { + state.lastCaughtUpTimestamp = OptionalLong.of(Math.max(fetchTimestamp, currentEndOffset.offset)); + } /*else { + state.lastCaughtUpTime = state.lastCaughtUpTime; + }*/ }); state.endOffset = Optional.of(endOffsetMetadata); @@ -290,6 +300,11 @@ private ReplicaState getReplicaState(int remoteNodeId) { return state; } + Long getReplicaLastFetchTimestamp(int remoteNodeId) { + ReplicaState replicaState = getReplicaState(remoteNodeId); + return replicaState.lastFetchTimestamp.orElse(-1L); + } + Map getVoterEndOffsets() { return getReplicaEndOffsets(voterStates); } @@ -323,12 +338,14 @@ private static class ReplicaState implements Comparable { final int nodeId; Optional endOffset; OptionalLong lastFetchTimestamp; + OptionalLong lastCaughtUpTimestamp; boolean hasAcknowledgedLeader; public ReplicaState(int nodeId, boolean hasAcknowledgedLeader) { this.nodeId = nodeId; this.endOffset = Optional.empty(); this.lastFetchTimestamp = OptionalLong.empty(); + this.lastCaughtUpTimestamp = OptionalLong.empty(); this.hasAcknowledgedLeader = hasAcknowledgedLeader; } @@ -338,6 +355,30 @@ void updateFetchTimestamp(long currentFetchTimeMs) { lastFetchTimestamp = OptionalLong.of(Math.max(lastFetchTimestamp.orElse(-1L), currentFetchTimeMs)); } + /** + * + * def updateFetchState( + * followerFetchOffsetMetadata: LogOffsetMetadata, + * followerStartOffset: Long, + * followerFetchTimeMs: Long, + * leaderEndOffset: Long + * ): Unit = { + * replicaState.updateAndGet { currentReplicaState => + * val lastCaughtUpTime = if (followerFetchOffsetMetadata.messageOffset >= leaderEndOffset) { + * math.max(currentReplicaState.lastCaughtUpTimeMs, followerFetchTimeMs) + * } else if (followerFetchOffsetMetadata.messageOffset >= currentReplicaState.lastFetchLeaderLogEndOffset) { + * math.max(currentReplicaState.lastCaughtUpTimeMs, currentReplicaState.lastFetchTimeMs) + * } else { + * currentReplicaState.lastCaughtUpTimeMs + * } + * + * + * @param currentFetchTimeMs + */ + public void updateLastCaughtUpTimeMs(long currentFetchTimeMs) { + + } + @Override public int compareTo(ReplicaState that) { if (this.endOffset.equals(that.endOffset)) diff --git a/raft/src/test/java/org/apache/kafka/raft/LeaderStateTest.java b/raft/src/test/java/org/apache/kafka/raft/LeaderStateTest.java index 5f9989d55e4c3..7060c0c28bae8 100644 --- a/raft/src/test/java/org/apache/kafka/raft/LeaderStateTest.java +++ b/raft/src/test/java/org/apache/kafka/raft/LeaderStateTest.java @@ -149,12 +149,12 @@ public void testUpdateHighWatermarkQuorumSizeTwo() { assertFalse(state.updateLocalState(0, new LogOffsetMetadata(13L))); assertEquals(singleton(otherNodeId), state.nonAcknowledgingVoters()); assertEquals(Optional.empty(), state.highWatermark()); - assertFalse(state.updateReplicaState(otherNodeId, 0, new LogOffsetMetadata(10L))); + assertFalse(state.updateReplicaState(otherNodeId, 0, new LogOffsetMetadata(10L), -1L)); assertEquals(emptySet(), state.nonAcknowledgingVoters()); assertEquals(Optional.empty(), state.highWatermark()); - assertTrue(state.updateReplicaState(otherNodeId, 0, new LogOffsetMetadata(11L))); + assertTrue(state.updateReplicaState(otherNodeId, 0, new LogOffsetMetadata(11L), -1L)); assertEquals(Optional.of(new LogOffsetMetadata(11L)), state.highWatermark()); - assertTrue(state.updateReplicaState(otherNodeId, 0, new LogOffsetMetadata(13L))); + assertTrue(state.updateReplicaState(otherNodeId, 0, new LogOffsetMetadata(13L), -1L)); assertEquals(Optional.of(new LogOffsetMetadata(13L)), state.highWatermark()); } @@ -166,19 +166,19 @@ public void testUpdateHighWatermarkQuorumSizeThree() { assertFalse(state.updateLocalState(0, new LogOffsetMetadata(15L))); assertEquals(mkSet(node1, node2), state.nonAcknowledgingVoters()); assertEquals(Optional.empty(), state.highWatermark()); - assertFalse(state.updateReplicaState(node1, 0, new LogOffsetMetadata(10L))); + assertFalse(state.updateReplicaState(node1, 0, new LogOffsetMetadata(10L), -1L)); assertEquals(singleton(node2), state.nonAcknowledgingVoters()); assertEquals(Optional.empty(), state.highWatermark()); - assertFalse(state.updateReplicaState(node2, 0, new LogOffsetMetadata(10L))); + assertFalse(state.updateReplicaState(node2, 0, new LogOffsetMetadata(10L), -1L)); assertEquals(emptySet(), state.nonAcknowledgingVoters()); assertEquals(Optional.empty(), state.highWatermark()); - assertTrue(state.updateReplicaState(node2, 0, new LogOffsetMetadata(15L))); + assertTrue(state.updateReplicaState(node2, 0, new LogOffsetMetadata(15L), -1L)); assertEquals(Optional.of(new LogOffsetMetadata(15L)), state.highWatermark()); assertFalse(state.updateLocalState(0, new LogOffsetMetadata(20L))); assertEquals(Optional.of(new LogOffsetMetadata(15L)), state.highWatermark()); - assertTrue(state.updateReplicaState(node1, 0, new LogOffsetMetadata(20L))); + assertTrue(state.updateReplicaState(node1, 0, new LogOffsetMetadata(20L), -1L)); assertEquals(Optional.of(new LogOffsetMetadata(20L)), state.highWatermark()); - assertFalse(state.updateReplicaState(node2, 0, new LogOffsetMetadata(20L))); + assertFalse(state.updateReplicaState(node2, 0, new LogOffsetMetadata(20L), -1L)); assertEquals(Optional.of(new LogOffsetMetadata(20L)), state.highWatermark()); } @@ -188,12 +188,12 @@ public void testNonMonotonicHighWatermarkUpdate() { int node1 = 1; LeaderState state = newLeaderState(mkSet(localId, node1), 0L); state.updateLocalState(time.milliseconds(), new LogOffsetMetadata(10L)); - state.updateReplicaState(node1, time.milliseconds(), new LogOffsetMetadata(10L)); + state.updateReplicaState(node1, time.milliseconds(), new LogOffsetMetadata(10L), -1L); assertEquals(Optional.of(new LogOffsetMetadata(10L)), state.highWatermark()); // Follower crashes and disk is lost. It fetches an earlier offset to rebuild state. // The leader will report an error in the logs, but will not let the high watermark rewind - assertFalse(state.updateReplicaState(node1, time.milliseconds(), new LogOffsetMetadata(5L))); + assertFalse(state.updateReplicaState(node1, time.milliseconds(), new LogOffsetMetadata(5L), -1L)); assertEquals(5L, state.getVoterEndOffsets().get(node1)); assertEquals(Optional.of(new LogOffsetMetadata(10L)), state.highWatermark()); } @@ -234,8 +234,8 @@ private LeaderState setUpLeaderAndFollowers(int follower1, LeaderState state = newLeaderState(mkSet(localId, follower1, follower2), leaderStartOffset); state.updateLocalState(0, new LogOffsetMetadata(leaderEndOffset)); assertEquals(Optional.empty(), state.highWatermark()); - state.updateReplicaState(follower1, 0, new LogOffsetMetadata(leaderStartOffset)); - state.updateReplicaState(follower2, 0, new LogOffsetMetadata(leaderEndOffset)); + state.updateReplicaState(follower1, 0, new LogOffsetMetadata(leaderStartOffset), -1L); + state.updateReplicaState(follower2, 0, new LogOffsetMetadata(leaderEndOffset), -1L); return state; } @@ -246,7 +246,7 @@ public void testGetObserverStatesWithObserver() { LeaderState state = newLeaderState(mkSet(localId), epochStartOffset); long timestamp = 20L; - assertFalse(state.updateReplicaState(observerId, timestamp, new LogOffsetMetadata(epochStartOffset))); + assertFalse(state.updateReplicaState(observerId, timestamp, new LogOffsetMetadata(epochStartOffset), -1L)); assertEquals(Collections.singletonMap(observerId, epochStartOffset), state.getObserverStates(timestamp)); } @@ -257,7 +257,7 @@ public void testNoOpForNegativeRemoteNodeId() { long epochStartOffset = 10L; LeaderState state = newLeaderState(mkSet(localId), epochStartOffset); - assertFalse(state.updateReplicaState(observerId, 0, new LogOffsetMetadata(epochStartOffset))); + assertFalse(state.updateReplicaState(observerId, 0, new LogOffsetMetadata(epochStartOffset), -1L)); assertEquals(Collections.emptyMap(), state.getObserverStates(10)); } @@ -269,7 +269,7 @@ public void testObserverStateExpiration() { long epochStartOffset = 10L; LeaderState state = newLeaderState(mkSet(localId), epochStartOffset); - state.updateReplicaState(observerId, time.milliseconds(), new LogOffsetMetadata(epochStartOffset)); + state.updateReplicaState(observerId, time.milliseconds(), new LogOffsetMetadata(epochStartOffset), -1L); assertEquals(singleton(observerId), state.getObserverStates(time.milliseconds()).keySet()); time.sleep(LeaderState.OBSERVER_SESSION_TIMEOUT_MS); diff --git a/raft/src/test/java/org/apache/kafka/raft/internals/KafkaRaftMetricsTest.java b/raft/src/test/java/org/apache/kafka/raft/internals/KafkaRaftMetricsTest.java index 0d64eac1cc2e2..cb9c03d8e4c21 100644 --- a/raft/src/test/java/org/apache/kafka/raft/internals/KafkaRaftMetricsTest.java +++ b/raft/src/test/java/org/apache/kafka/raft/internals/KafkaRaftMetricsTest.java @@ -103,7 +103,7 @@ public void shouldRecordVoterQuorumState() throws IOException { assertEquals((double) -1L, getMetric(metrics, "high-watermark").metricValue()); state.leaderStateOrThrow().updateLocalState(0, new LogOffsetMetadata(5L)); - state.leaderStateOrThrow().updateReplicaState(1, 0, new LogOffsetMetadata(5L)); + state.leaderStateOrThrow().updateReplicaState(1, 0, new LogOffsetMetadata(5L), -1L); assertEquals((double) 5L, getMetric(metrics, "high-watermark").metricValue()); state.transitionToFollower(2, 1);