diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index a2b07913356db..b1972b71a1493 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -321,13 +321,13 @@
+
-
consumerGr
subscriptionMetadata = group.computeSubscriptionMetadata(
member,
updatedMember,
- metadataImage.topics()
+ metadataImage.topics(),
+ metadataImage.cluster()
);
if (!subscriptionMetadata.equals(group.subscriptionMetadata())) {
@@ -927,7 +928,8 @@ private List consumerGroupFenceMember(
Map subscriptionMetadata = group.computeSubscriptionMetadata(
member,
null,
- metadataImage.topics()
+ metadataImage.topics(),
+ metadataImage.cluster()
);
if (!subscriptionMetadata.equals(group.subscriptionMetadata())) {
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/RecordHelpers.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/RecordHelpers.java
index 14a55b87324c7..b5ef473dd4923 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/RecordHelpers.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/RecordHelpers.java
@@ -131,13 +131,21 @@ public static Record newGroupSubscriptionMetadataRecord(
Map newSubscriptionMetadata
) {
ConsumerGroupPartitionMetadataValue value = new ConsumerGroupPartitionMetadataValue();
- newSubscriptionMetadata.forEach((topicName, topicMetadata) ->
+ newSubscriptionMetadata.forEach((topicName, topicMetadata) -> {
+ List partitionMetadata = new ArrayList<>();
+ topicMetadata.partitionRacks().forEach((partition, racks) ->
+ partitionMetadata.add(new ConsumerGroupPartitionMetadataValue.PartitionMetadata()
+ .setPartition(partition)
+ .setRacks(new ArrayList<>(racks))
+ )
+ );
value.topics().add(new ConsumerGroupPartitionMetadataValue.TopicMetadata()
.setTopicId(topicMetadata.id())
.setTopicName(topicMetadata.name())
.setNumPartitions(topicMetadata.numPartitions())
- )
- );
+ .setPartitionMetadata(partitionMetadata)
+ );
+ });
return new Record(
new ApiMessageAndVersion(
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentSpec.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentSpec.java
index bf775439077ea..fbc2c20626924 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentSpec.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentSpec.java
@@ -16,8 +16,6 @@
*/
package org.apache.kafka.coordinator.group.assignor;
-import org.apache.kafka.common.Uuid;
-
import java.util.Map;
import java.util.Objects;
@@ -26,59 +24,39 @@
*/
public class AssignmentSpec {
/**
- * The members keyed by member id.
+ * The members keyed by member Id.
*/
private final Map members;
- /**
- * The topics' metadata keyed by topic id.
- */
- private final Map topics;
-
public AssignmentSpec(
- Map members,
- Map topics
+ Map members
) {
Objects.requireNonNull(members);
- Objects.requireNonNull(topics);
this.members = members;
- this.topics = topics;
}
/**
- * @return Member metadata keyed by member Ids.
+ * @return Member metadata keyed by member Id.
*/
public Map members() {
return members;
}
- /**
- * @return Topic metadata keyed by topic Ids.
- */
- public Map topics() {
- return topics;
- }
-
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
AssignmentSpec that = (AssignmentSpec) o;
- if (!members.equals(that.members)) return false;
- return topics.equals(that.topics);
+ return members.equals(that.members);
}
@Override
public int hashCode() {
- int result = members.hashCode();
- result = 31 * result + topics.hashCode();
- return result;
+ return members.hashCode();
}
@Override
public String toString() {
- return "AssignmentSpec(members=" + members +
- ", topics=" + topics +
- ')';
+ return "AssignmentSpec(members=" + members + ')';
}
}
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentTopicMetadata.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentTopicMetadata.java
deleted file mode 100644
index e5a96e79b2ca7..0000000000000
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/AssignmentTopicMetadata.java
+++ /dev/null
@@ -1,59 +0,0 @@
-/*
- * 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.coordinator.group.assignor;
-
-/**
- * Metadata of a topic.
- */
-public class AssignmentTopicMetadata {
-
- /**
- * The number of partitions.
- */
- private final int numPartitions;
-
- public AssignmentTopicMetadata(
- int numPartitions
- ) {
- this.numPartitions = numPartitions;
- }
-
- /**
- * @return The number of partitions present for the topic.
- */
- public int numPartitions() {
- return numPartitions;
- }
-
- @Override
- public boolean equals(Object o) {
- if (this == o) return true;
- if (o == null || getClass() != o.getClass()) return false;
- AssignmentTopicMetadata that = (AssignmentTopicMetadata) o;
- return numPartitions == that.numPartitions;
- }
-
- @Override
- public int hashCode() {
- return numPartitions;
- }
-
- @Override
- public String toString() {
- return "AssignmentTopicMetadata(numPartitions=" + numPartitions + ')';
- }
-}
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/PartitionAssignor.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/PartitionAssignor.java
index 43410fc8aae40..7fd74889fef8b 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/PartitionAssignor.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/PartitionAssignor.java
@@ -36,8 +36,9 @@ public interface PartitionAssignor {
* Perform the group assignment given the current members and
* topic metadata.
*
- * @param assignmentSpec The assignment spec.
+ * @param assignmentSpec The member assignment spec.
+ * @param subscribedTopicDescriber The topic and cluster metadata describer {@link SubscribedTopicDescriber}.
* @return The new assignment for the group.
*/
- GroupAssignment assign(AssignmentSpec assignmentSpec) throws PartitionAssignorException;
+ GroupAssignment assign(AssignmentSpec assignmentSpec, SubscribedTopicDescriber subscribedTopicDescriber) throws PartitionAssignorException;
}
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/RangeAssignor.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/RangeAssignor.java
index b5316ffbbdbf2..4380c0fd8a3ff 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/RangeAssignor.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/RangeAssignor.java
@@ -78,23 +78,19 @@ public MemberWithRemainingAssignments(String memberId, int remaining) {
/**
* @return Map of topic ids to a list of members subscribed to them.
*/
- private Map> membersPerTopic(final AssignmentSpec assignmentSpec) {
+ private Map> membersPerTopic(final AssignmentSpec assignmentSpec, final SubscribedTopicDescriber subscribedTopicDescriber) {
Map> membersPerTopic = new HashMap<>();
Map membersData = assignmentSpec.members();
membersData.forEach((memberId, memberMetadata) -> {
Collection topics = memberMetadata.subscribedTopicIds();
for (Uuid topicId : topics) {
- // Only topics that are present in both the subscribed topics list and the topic metadata should be
- // considered for assignment.
- if (assignmentSpec.topics().containsKey(topicId)) {
- membersPerTopic
- .computeIfAbsent(topicId, k -> new ArrayList<>())
- .add(memberId);
- } else {
- throw new PartitionAssignorException("Member " + memberId + " subscribed to topic " +
- topicId + " which doesn't exist in the topic metadata");
+ if (subscribedTopicDescriber.numPartitions(topicId) == -1) {
+ throw new PartitionAssignorException("Member is subscribed to a non-existent topic");
}
+ membersPerTopic
+ .computeIfAbsent(topicId, k -> new ArrayList<>())
+ .add(memberId);
}
});
@@ -118,14 +114,14 @@ private Map> membersPerTopic(final AssignmentSpec assignmentS
*
*/
@Override
- public GroupAssignment assign(final AssignmentSpec assignmentSpec) throws PartitionAssignorException {
+ public GroupAssignment assign(final AssignmentSpec assignmentSpec, final SubscribedTopicDescriber subscribedTopicDescriber) throws PartitionAssignorException {
Map newAssignment = new HashMap<>();
// Step 1
- Map> membersPerTopic = membersPerTopic(assignmentSpec);
+ Map> membersPerTopic = membersPerTopic(assignmentSpec, subscribedTopicDescriber);
membersPerTopic.forEach((topicId, membersForTopic) -> {
- int numPartitionsForTopic = assignmentSpec.topics().get(topicId).numPartitions();
+ int numPartitionsForTopic = subscribedTopicDescriber.numPartitions(topicId);
int minRequiredQuota = numPartitionsForTopic / membersForTopic.size();
// Each member can get only ONE extra partition per topic after receiving the minimum quota.
int numMembersWithExtraPartition = numPartitionsForTopic % membersForTopic.size();
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/SubscribedTopicDescriber.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/SubscribedTopicDescriber.java
new file mode 100644
index 0000000000000..80fd3b737b5f7
--- /dev/null
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/assignor/SubscribedTopicDescriber.java
@@ -0,0 +1,51 @@
+/*
+ * 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.coordinator.group.assignor;
+
+import org.apache.kafka.common.Uuid;
+import org.apache.kafka.common.annotation.InterfaceStability;
+
+import java.util.Set;
+
+/**
+ * The subscribed topic describer is used by the {@link PartitionAssignor}
+ * to obtain topic and partition metadata of subscribed topics.
+ *
+ * The interface is kept in an internal module until KIP-848 is fully
+ * implemented and ready to be released.
+ */
+@InterfaceStability.Unstable
+public interface SubscribedTopicDescriber {
+ /**
+ * Number of partitions for the given topic Id.
+ *
+ * @param topicId Uuid corresponding to the topic.
+ * @return The number of partitions corresponding to the given topicId.
+ * If the topicId doesn't exist return -1;
+ */
+ int numPartitions(Uuid topicId);
+
+ /**
+ * Returns all the racks associated with the replicas for the given partition.
+ *
+ * @param topicId Uuid corresponding to the partition's topic.
+ * @param partition Partition number within topic.
+ * @return The set of racks corresponding to the replicas of the topics partition.
+ * If the topicId doesn't exist return an empty set;
+ */
+ Set racksForPartition(Uuid topicId, int partition);
+}
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java
index 6f682d4e82f87..5df816c5d9cb6 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroup.java
@@ -19,6 +19,7 @@
import org.apache.kafka.common.Uuid;
import org.apache.kafka.common.errors.UnknownMemberIdException;
import org.apache.kafka.coordinator.group.Group;
+import org.apache.kafka.image.ClusterImage;
import org.apache.kafka.image.TopicImage;
import org.apache.kafka.image.TopicsImage;
import org.apache.kafka.timeline.SnapshotRegistry;
@@ -28,6 +29,7 @@
import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
@@ -434,7 +436,8 @@ public void setSubscriptionMetadata(
public Map computeSubscriptionMetadata(
ConsumerGroupMember oldMember,
ConsumerGroupMember newMember,
- TopicsImage topicsImage
+ TopicsImage topicsImage,
+ ClusterImage clusterImage
) {
// Copy and update the current subscriptions.
Map subscribedTopicNames = new HashMap<>(this.subscribedTopicNames);
@@ -442,14 +445,30 @@ public Map computeSubscriptionMetadata(
// Create the topic metadata for each subscribed topic.
Map newSubscriptionMetadata = new HashMap<>(subscribedTopicNames.size());
+
subscribedTopicNames.forEach((topicName, count) -> {
TopicImage topicImage = topicsImage.getTopic(topicName);
if (topicImage != null) {
+ Map> partitionRacks = new HashMap<>();
+
+ topicImage.partitions().forEach((partition, partitionRegistration) -> {
+ Set racks = new HashSet<>();
+ for (int replica : partitionRegistration.replicas) {
+ Optional rackOptional = clusterImage.broker(replica).rack();
+ // Only add rack if it is available for the broker/replica.
+ rackOptional.ifPresent(racks::add);
+ }
+ // If no racks are available for any replica of this partition, store an empty map.
+ if (!racks.isEmpty())
+ partitionRacks.put(partition, racks);
+ });
+
newSubscriptionMetadata.put(topicName, new TopicMetadata(
topicImage.id(),
topicImage.name(),
- topicImage.partitions().size()
- ));
+ topicImage.partitions().size(),
+ partitionRacks)
+ );
}
});
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/SubscribedTopicMetadata.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/SubscribedTopicMetadata.java
new file mode 100644
index 0000000000000..095277012b6ad
--- /dev/null
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/SubscribedTopicMetadata.java
@@ -0,0 +1,89 @@
+/*
+ * 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.coordinator.group.consumer;
+
+import org.apache.kafka.common.Uuid;
+import org.apache.kafka.coordinator.group.assignor.PartitionAssignor;
+import org.apache.kafka.coordinator.group.assignor.SubscribedTopicDescriber;
+
+import java.util.Collections;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * The subscribed topic metadata class is used by the {@link PartitionAssignor}
+ * to obtain topic and partition metadata of subscribed topics.
+ */
+public class SubscribedTopicMetadata implements SubscribedTopicDescriber {
+
+ /**
+ * The topic IDs are mapped to their corresponding {@link TopicMetadata}
+ * object, which contains topic and partition metadata.
+ */
+ Map topicMetadata;
+
+ public SubscribedTopicMetadata(Map topicMetadata) {
+ this.topicMetadata = topicMetadata;
+ }
+
+ /**
+ * Number of partitions available for the given topic Id.
+ *
+ * @param topicId Uuid corresponding to the topic.
+ * @return The number of partitions corresponding to the given topicId.
+ * If the topicId doesn't exist return -1;
+ */
+ @Override
+ public int numPartitions(Uuid topicId) {
+ return this.topicMetadata.containsKey(topicId) ? this.topicMetadata.get(topicId).numPartitions() : -1;
+ }
+
+ /**
+ * Returns all the racks associated with the replicas for the given partition.
+ *
+ * @param topicId Uuid corresponding to the partition's topic.
+ * @param partition Partition number within topic.
+ * @return The set of racks corresponding to the replicas of the topics partition.
+ * If the topicId doesn't exist return an empty set;
+ */
+ @Override
+ public Set racksForPartition(Uuid topicId, int partition) {
+ return this.topicMetadata.containsKey(topicId) ?
+ this.topicMetadata.get(topicId).partitionRacks().get(partition) :
+ Collections.emptySet();
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (o == null || getClass() != o.getClass()) return false;
+ SubscribedTopicMetadata that = (SubscribedTopicMetadata) o;
+ return topicMetadata.equals(that.topicMetadata);
+ }
+
+ @Override
+ public int hashCode() {
+ return topicMetadata.hashCode();
+ }
+
+ @Override
+ public String toString() {
+ return "SubscribedTopicMetadata(" +
+ "topicMetadata=" + topicMetadata +
+ ')';
+ }
+}
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilder.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilder.java
index 02b120db1ef32..2026efaa82c9d 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilder.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilder.java
@@ -20,7 +20,6 @@
import org.apache.kafka.coordinator.group.Record;
import org.apache.kafka.coordinator.group.assignor.AssignmentMemberSpec;
import org.apache.kafka.coordinator.group.assignor.AssignmentSpec;
-import org.apache.kafka.coordinator.group.assignor.AssignmentTopicMetadata;
import org.apache.kafka.coordinator.group.assignor.GroupAssignment;
import org.apache.kafka.coordinator.group.assignor.MemberAssignment;
import org.apache.kafka.coordinator.group.assignor.PartitionAssignor;
@@ -244,16 +243,19 @@ public TargetAssignmentResult build() throws PartitionAssignorException {
});
// Prepare the topic metadata.
- Map topics = new HashMap<>();
+ Map topicMetadataMap = new HashMap<>();
subscriptionMetadata.forEach((topicName, topicMetadata) ->
- topics.put(topicMetadata.id(), new AssignmentTopicMetadata(topicMetadata.numPartitions()))
+ topicMetadataMap.put(
+ topicMetadata.id(),
+ topicMetadata
+ )
);
// Compute the assignment.
- GroupAssignment newGroupAssignment = assignor.assign(new AssignmentSpec(
- Collections.unmodifiableMap(memberSpecs),
- Collections.unmodifiableMap(topics)
- ));
+ GroupAssignment newGroupAssignment = assignor.assign(
+ new AssignmentSpec(Collections.unmodifiableMap(memberSpecs)),
+ new SubscribedTopicMetadata(topicMetadataMap)
+ );
// Compute delta from previous to new target assignment and create the
// relevant records.
diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TopicMetadata.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TopicMetadata.java
index 395d48d95dffb..5c1d4b1c5bfb4 100644
--- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TopicMetadata.java
+++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/consumer/TopicMetadata.java
@@ -19,7 +19,12 @@
import org.apache.kafka.common.Uuid;
import org.apache.kafka.coordinator.group.generated.ConsumerGroupPartitionMetadataValue;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
import java.util.Objects;
+import java.util.Set;
/**
* Immutable topic metadata.
@@ -40,23 +45,31 @@ public class TopicMetadata {
*/
private final int numPartitions;
+ /**
+ * Map of every partition to a set of its rackIds.
+ * If the rack information is unavailable, this is an empty map.
+ */
+ private final Map> partitionRacks;
+
public TopicMetadata(
Uuid id,
String name,
- int numPartitions
+ int numPartitions,
+ Map> partitionRacks
) {
- this.id = Objects.requireNonNull(id);
- if (Uuid.ZERO_UUID.equals(id)) {
- throw new IllegalArgumentException("Topic id cannot be ZERO_UUID.");
- }
- this.name = Objects.requireNonNull(name);
- if (name.isEmpty()) {
- throw new IllegalArgumentException("Topic name cannot be empty.");
- }
- this.numPartitions = numPartitions;
- if (numPartitions < 0) {
- throw new IllegalArgumentException("Number of partitions cannot be negative.");
- }
+ this.id = Objects.requireNonNull(id);
+ this.partitionRacks = Objects.requireNonNull(partitionRacks);
+ if (Uuid.ZERO_UUID.equals(id)) {
+ throw new IllegalArgumentException("Topic id cannot be ZERO_UUID.");
+ }
+ this.name = Objects.requireNonNull(name);
+ if (name.isEmpty()) {
+ throw new IllegalArgumentException("Topic name cannot be empty.");
+ }
+ this.numPartitions = numPartitions;
+ if (numPartitions < 0) {
+ throw new IllegalArgumentException("Number of partitions cannot be negative.");
+ }
}
/**
@@ -80,6 +93,14 @@ public int numPartitions() {
return this.numPartitions;
}
+ /**
+ * @return Every partition mapped to the set of corresponding rack Ids of its replicas.
+ * An empty map is returned if no rack information is available.
+ */
+ public Map> partitionRacks() {
+ return this.partitionRacks;
+ }
+
@Override
public boolean equals(Object o) {
if (this == o) return true;
@@ -89,7 +110,8 @@ public boolean equals(Object o) {
if (!id.equals(that.id)) return false;
if (!name.equals(that.name)) return false;
- return numPartitions == that.numPartitions;
+ if (numPartitions != that.numPartitions) return false;
+ return partitionRacks.equals(that.partitionRacks);
}
@Override
@@ -97,6 +119,7 @@ public int hashCode() {
int result = id.hashCode();
result = 31 * result + name.hashCode();
result = 31 * result + numPartitions;
+ result = 31 * result + partitionRacks.hashCode();
return result;
}
@@ -106,16 +129,23 @@ public String toString() {
"id=" + id +
", name=" + name +
", numPartitions=" + numPartitions +
+ ", partitionRacks=" + partitionRacks +
')';
}
public static TopicMetadata fromRecord(
ConsumerGroupPartitionMetadataValue.TopicMetadata record
) {
+ // Converting the data type from a list stored in the record to a map.
+ Map> partitionRacks = new HashMap<>(record.partitionMetadata().size());
+ for (ConsumerGroupPartitionMetadataValue.PartitionMetadata partitionMetadata : record.partitionMetadata()) {
+ partitionRacks.put(partitionMetadata.partition(), Collections.unmodifiableSet(new HashSet<>(partitionMetadata.racks())));
+ }
+
return new TopicMetadata(
record.topicId(),
record.topicName(),
- record.numPartitions()
- );
+ record.numPartitions(),
+ partitionRacks);
}
}
diff --git a/group-coordinator/src/main/resources/common/message/ConsumerGroupPartitionMetadataValue.json b/group-coordinator/src/main/resources/common/message/ConsumerGroupPartitionMetadataValue.json
index 81b5f5225e7bc..260792df1907a 100644
--- a/group-coordinator/src/main/resources/common/message/ConsumerGroupPartitionMetadataValue.json
+++ b/group-coordinator/src/main/resources/common/message/ConsumerGroupPartitionMetadataValue.json
@@ -29,7 +29,14 @@
{ "name": "TopicName", "versions": "0+", "type": "string",
"about": "The topic name." },
{ "name": "NumPartitions", "versions": "0+", "type": "int32",
- "about": "The number of partitions of the topic." }
+ "about": "The number of partitions of the topic." },
+ { "name": "PartitionMetadata", "versions": "0+", "type": "[]PartitionMetadata",
+ "about": "Partitions mapped to a set of racks.", "fields": [
+ { "name": "Partition", "versions": "0+", "type": "int32",
+ "about": "The partition number." },
+ { "name": "Racks", "versions": "0+", "type": "[]string",
+ "about": "The set of racks that the partition is mapped to." }
+ ]}
]}
]
}
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
index af4cb709b2801..bbb5d538faf63 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
@@ -33,6 +33,7 @@
import org.apache.kafka.common.message.JoinGroupRequestData;
import org.apache.kafka.common.message.JoinGroupResponseData;
import org.apache.kafka.common.metadata.PartitionRecord;
+import org.apache.kafka.common.metadata.RegisterBrokerRecord;
import org.apache.kafka.common.metadata.RemoveTopicRecord;
import org.apache.kafka.common.metadata.TopicRecord;
import org.apache.kafka.common.network.ClientInformation;
@@ -53,6 +54,7 @@
import org.apache.kafka.coordinator.group.assignor.MemberAssignment;
import org.apache.kafka.coordinator.group.assignor.PartitionAssignor;
import org.apache.kafka.coordinator.group.assignor.PartitionAssignorException;
+import org.apache.kafka.coordinator.group.assignor.SubscribedTopicDescriber;
import org.apache.kafka.coordinator.group.consumer.Assignment;
import org.apache.kafka.coordinator.group.consumer.ConsumerGroup;
import org.apache.kafka.coordinator.group.consumer.ConsumerGroupMember;
@@ -115,6 +117,7 @@
import static org.apache.kafka.coordinator.group.GroupMetadataManager.consumerGroupSessionTimeoutKey;
import static org.apache.kafka.coordinator.group.GroupMetadataManager.EMPTY_RESULT;
import static org.apache.kafka.coordinator.group.GroupMetadataManager.genericGroupHeartbeatKey;
+import static org.apache.kafka.coordinator.group.RecordHelpersTest.mkMapOfPartitionRacks;
import static org.apache.kafka.coordinator.group.generic.GenericGroupState.COMPLETING_REBALANCE;
import static org.apache.kafka.coordinator.group.generic.GenericGroupState.DEAD;
import static org.apache.kafka.coordinator.group.generic.GenericGroupState.EMPTY;
@@ -152,7 +155,7 @@ public String name() {
}
@Override
- public GroupAssignment assign(AssignmentSpec assignmentSpec) throws PartitionAssignorException {
+ public GroupAssignment assign(AssignmentSpec assignmentSpec, SubscribedTopicDescriber subscribedTopicDescriber) throws PartitionAssignorException {
return prepareGroupAssignment;
}
}
@@ -174,6 +177,14 @@ public MetadataImageBuilder addTopic(
return this;
}
+ public MetadataImageBuilder addRack(
+ int brokerId,
+ String rack
+ ) {
+ delta.replay(new RegisterBrokerRecord().setBrokerId(brokerId).setRack(rack));
+ return this;
+ }
+
public MetadataImage build() {
return delta.apply(MetadataProvenance.EMPTY);
}
@@ -231,8 +242,8 @@ public List build(TopicsImage topicsImage) {
subscriptionMetadata.put(topicName, new TopicMetadata(
topicImage.id(),
topicImage.name(),
- topicImage.partitions().size()
- ));
+ topicImage.partitions().size(),
+ Collections.emptyMap()));
}
});
});
@@ -1042,8 +1053,8 @@ public void testMemberJoinsEmptyConsumerGroup() {
List expectedRecords = Arrays.asList(
RecordHelpers.newMemberSubscriptionRecord(groupId, expectedMember),
RecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {{
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6));
- put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 3));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6, mkMapOfPartitionRacks(6)));
+ put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 3, mkMapOfPartitionRacks(3)));
}}),
RecordHelpers.newGroupEpochRecord(groupId, 1),
RecordHelpers.newTargetAssignmentRecord(groupId, memberId, mkAssignment(
@@ -1082,7 +1093,7 @@ public void testUpdatingSubscriptionTriggersNewTargetAssignment() {
.setTargetMemberEpoch(10)
.setClientId("client")
.setClientHost("localhost/127.0.0.1")
- .setSubscribedTopicNames(Arrays.asList("foo"))
+ .setSubscribedTopicNames(Collections.singletonList("foo"))
.setServerAssignorName("range")
.setAssignedPartitions(mkAssignment(
mkTopicAssignment(fooTopicId, 0, 1, 2, 3, 4, 5)))
@@ -1140,8 +1151,8 @@ public void testUpdatingSubscriptionTriggersNewTargetAssignment() {
RecordHelpers.newMemberSubscriptionRecord(groupId, expectedMember),
RecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
{
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6));
- put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 3));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6, mkMapOfPartitionRacks(6)));
+ put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 3, mkMapOfPartitionRacks(3)));
}
}),
RecordHelpers.newGroupEpochRecord(groupId, 11),
@@ -1254,7 +1265,7 @@ public void testNewJoiningMemberTriggersNewTargetAssignment() {
.setPartitions(Arrays.asList(4, 5)),
new ConsumerGroupHeartbeatResponseData.TopicPartitions()
.setTopicId(barTopicId)
- .setPartitions(Arrays.asList(2))
+ .setPartitions(Collections.singletonList(2))
))),
result.response()
);
@@ -1380,8 +1391,8 @@ public void testLeavingMemberBumpsGroupEpoch() {
// Subscription metadata is recomputed because zar is no longer there.
RecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
{
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6));
- put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 3));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6, mkMapOfPartitionRacks(6)));
+ put(barTopicName, new TopicMetadata(barTopicId, barTopicName, 3, mkMapOfPartitionRacks(3)));
}
}),
RecordHelpers.newGroupEpochRecord(groupId, 11)
@@ -2215,7 +2226,7 @@ public void testPartitionAssignorExceptionOnRegularHeartbeat() {
PartitionAssignor assignor = mock(PartitionAssignor.class);
when(assignor.name()).thenReturn("range");
- when(assignor.assign(any())).thenThrow(new PartitionAssignorException("Assignment failed."));
+ when(assignor.assign(any(), any())).thenThrow(new PartitionAssignorException("Assignment failed."));
GroupMetadataManagerTestContext context = new GroupMetadataManagerTestContext.Builder()
.withAssignors(Collections.singletonList(assignor))
@@ -2276,7 +2287,7 @@ public void testSubscriptionMetadataRefreshedAfterGroupIsLoaded() {
{
// foo only has 3 partitions stored in the metadata but foo has
// 6 partitions the metadata image.
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 3));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 3, mkMapOfPartitionRacks(3)));
}
}))
.build();
@@ -2330,7 +2341,7 @@ public void testSubscriptionMetadataRefreshedAfterGroupIsLoaded() {
List expectedRecords = Arrays.asList(
RecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
{
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6, mkMapOfPartitionRacks(6)));
}
}),
RecordHelpers.newGroupEpochRecord(groupId, 11),
@@ -2386,7 +2397,7 @@ public void testSubscriptionMetadataRefreshedAgainAfterWriteFailure() {
{
// foo only has 3 partitions stored in the metadata but foo has
// 6 partitions the metadata image.
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 3));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 3, mkMapOfPartitionRacks(3)));
}
}))
.build();
@@ -2458,7 +2469,7 @@ public void testSubscriptionMetadataRefreshedAgainAfterWriteFailure() {
List expectedRecords = Arrays.asList(
RecordHelpers.newGroupSubscriptionMetadataRecord(groupId, new HashMap() {
{
- put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6));
+ put(fooTopicName, new TopicMetadata(fooTopicId, fooTopicName, 6, mkMapOfPartitionRacks(6)));
}
}),
RecordHelpers.newGroupEpochRecord(groupId, 11),
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/RecordHelpersTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/RecordHelpersTest.java
index accda00808a53..0480a292e88e8 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/RecordHelpersTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/RecordHelpersTest.java
@@ -58,6 +58,7 @@
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -160,16 +161,17 @@ public void testNewGroupSubscriptionMetadataRecord() {
Uuid fooTopicId = Uuid.randomUuid();
Uuid barTopicId = Uuid.randomUuid();
Map subscriptionMetadata = new LinkedHashMap<>();
+
subscriptionMetadata.put("foo", new TopicMetadata(
fooTopicId,
"foo",
- 10
- ));
+ 10,
+ mkMapOfPartitionRacks(10)));
subscriptionMetadata.put("bar", new TopicMetadata(
barTopicId,
"bar",
- 20
- ));
+ 20,
+ mkMapOfPartitionRacks(20)));
Record expectedRecord = new Record(
new ApiMessageAndVersion(
@@ -183,11 +185,13 @@ public void testNewGroupSubscriptionMetadataRecord() {
new ConsumerGroupPartitionMetadataValue.TopicMetadata()
.setTopicId(fooTopicId)
.setTopicName("foo")
- .setNumPartitions(10),
+ .setNumPartitions(10)
+ .setPartitionMetadata(mkListOfPartitionRacks(10)),
new ConsumerGroupPartitionMetadataValue.TopicMetadata()
.setTopicId(barTopicId)
.setTopicName("bar")
- .setNumPartitions(20))),
+ .setNumPartitions(20)
+ .setPartitionMetadata(mkListOfPartitionRacks(20)))),
(short) 0));
assertEquals(expectedRecord, newGroupSubscriptionMetadataRecord(
@@ -612,7 +616,7 @@ public void testNewGroupMetadataRecordThrowsWhenEmptyAssignment() {
MetadataVersion.IBP_3_5_IV2
));
}
-
+
@ParameterizedTest
@MethodSource("metadataToExpectedGroupMetadataValue")
public void testEmptyGroupMetadataRecord(
@@ -620,7 +624,6 @@ public void testEmptyGroupMetadataRecord(
short expectedGroupMetadataValueVersion
) {
Time time = new MockTime();
-
List expectedMembers = Collections.emptyList();
Record expectedRecord = new Record(
@@ -756,4 +759,24 @@ public void testNewOffsetCommitTombstoneRecord() {
Record record = RecordHelpers.newOffsetCommitTombstoneRecord("group-id", "foo", 1);
assertEquals(expectedRecord, record);
}
+
+ public static List mkListOfPartitionRacks(int numPartitions) {
+ List partitionRacks = new ArrayList<>(numPartitions);
+ for (int i = 0; i < numPartitions; i++) {
+ partitionRacks.add(
+ new ConsumerGroupPartitionMetadataValue.PartitionMetadata()
+ .setPartition(i)
+ .setRacks(Collections.emptyList())
+ );
+ }
+ return partitionRacks;
+ }
+
+ public static Map> mkMapOfPartitionRacks(int numPartitions) {
+ Map> partitionRacks = new HashMap<>(numPartitions);
+ for (int i = 0; i < numPartitions ; i++) {
+ partitionRacks.put(i, Collections.emptySet());
+ }
+ return partitionRacks;
+ }
}
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/RangeAssignorTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/RangeAssignorTest.java
index 91f6385f104e4..24e820169cd4d 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/RangeAssignorTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/assignor/RangeAssignorTest.java
@@ -18,6 +18,8 @@
import org.apache.kafka.common.Uuid;
+import org.apache.kafka.coordinator.group.consumer.SubscribedTopicMetadata;
+import org.apache.kafka.coordinator.group.consumer.TopicMetadata;
import org.junit.jupiter.api.Test;
import java.util.Arrays;
@@ -37,15 +39,29 @@
public class RangeAssignorTest {
private final RangeAssignor assignor = new RangeAssignor();
private final Uuid topic1Uuid = Uuid.randomUuid();
+ private final String topic1Name = "topic1";
private final Uuid topic2Uuid = Uuid.randomUuid();
+ private final String topic2Name = "topic2";
private final Uuid topic3Uuid = Uuid.randomUuid();
+ private final String topic3Name = "topic3";
private final String consumerA = "A";
private final String consumerB = "B";
private final String consumerC = "C";
@Test
public void testOneConsumerNoTopic() {
- Map topics = Collections.singletonMap(topic1Uuid, new AssignmentTopicMetadata(3));
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(
+ Collections.singletonMap(
+ topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ )
+ );
+
Map members = Collections.singletonMap(
consumerA,
new AssignmentMemberSpec(
@@ -56,15 +72,26 @@ public void testOneConsumerNoTopic() {
)
);
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment groupAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment groupAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
assertEquals(Collections.emptyMap(), groupAssignment.members());
}
@Test
public void testOneConsumerSubscribedToNonExistentTopic() {
- Map topics = Collections.singletonMap(topic1Uuid, new AssignmentTopicMetadata(3));
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(
+ Collections.singletonMap(
+ topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ )
+ );
+
Map members = Collections.singletonMap(
consumerA,
new AssignmentMemberSpec(
@@ -75,17 +102,34 @@ public void testOneConsumerSubscribedToNonExistentTopic() {
)
);
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
assertThrows(PartitionAssignorException.class,
- () -> assignor.assign(assignmentSpec));
+ () -> assignor.assign(assignmentSpec, subscribedTopicMetadata));
}
@Test
public void testFirstAssignmentTwoConsumersTwoTopicsSameSubscriptions() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic3Uuid, new AssignmentTopicMetadata(2));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic3Uuid,
+ new TopicMetadata(
+ topic3Uuid,
+ topic3Name,
+ 2,
+ createPartitionMetadata(2)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -103,8 +147,8 @@ public void testFirstAssignmentTwoConsumersTwoTopicsSameSubscriptions() {
Collections.emptyMap()
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -123,10 +167,35 @@ public void testFirstAssignmentTwoConsumersTwoTopicsSameSubscriptions() {
@Test
public void testFirstAssignmentThreeConsumersThreeTopicsDifferentSubscriptions() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic3Uuid, new AssignmentTopicMetadata(2));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic3Uuid,
+ new TopicMetadata(
+ topic3Uuid,
+ topic3Name,
+ 2,
+ createPartitionMetadata(2)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -151,8 +220,8 @@ public void testFirstAssignmentThreeConsumersThreeTopicsDifferentSubscriptions()
Collections.emptyMap()
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -175,9 +244,26 @@ public void testFirstAssignmentThreeConsumersThreeTopicsDifferentSubscriptions()
@Test
public void testFirstAssignmentNumConsumersGreaterThanNumPartitions() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic3Uuid, new AssignmentTopicMetadata(2));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic3Uuid,
+ new TopicMetadata(
+ topic3Uuid,
+ topic3Name,
+ 2,
+ createPartitionMetadata(2)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -202,8 +288,8 @@ public void testFirstAssignmentNumConsumersGreaterThanNumPartitions() {
Collections.emptyMap()
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
// Topic 3 has 2 partitions but three consumers subscribed to it - one of them will not get a partition.
@@ -226,9 +312,26 @@ public void testFirstAssignmentNumConsumersGreaterThanNumPartitions() {
@Test
public void testReassignmentNumConsumersGreaterThanNumPartitionsWhenOneConsumerAdded() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(2));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(2));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 2,
+ createPartitionMetadata(2)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 2,
+ createPartitionMetadata(2)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -264,8 +367,8 @@ public void testReassignmentNumConsumersGreaterThanNumPartitionsWhenOneConsumerA
Collections.emptyMap()
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -287,9 +390,26 @@ public void testReassignmentNumConsumersGreaterThanNumPartitionsWhenOneConsumerA
@Test
public void testReassignmentWhenOnePartitionAddedForTwoConsumersTwoTopics() {
// Simulating adding a partition - originally T1 -> 3 Partitions and T2 -> 3 Partitions
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(4));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(4));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 4,
+ createPartitionMetadata(4)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 4,
+ createPartitionMetadata(4)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -317,8 +437,8 @@ public void testReassignmentWhenOnePartitionAddedForTwoConsumersTwoTopics() {
currentAssignmentForB
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -337,9 +457,26 @@ public void testReassignmentWhenOnePartitionAddedForTwoConsumersTwoTopics() {
@Test
public void testReassignmentWhenOneConsumerAddedAfterInitialAssignmentWithTwoConsumersTwoTopics() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(3));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -375,8 +512,8 @@ public void testReassignmentWhenOneConsumerAddedAfterInitialAssignmentWithTwoCon
Collections.emptyMap()
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -400,10 +537,27 @@ public void testReassignmentWhenOneConsumerAddedAfterInitialAssignmentWithTwoCon
@Test
public void testReassignmentWhenOneConsumerAddedAndOnePartitionAfterInitialAssignmentWithTwoConsumersTwoTopics() {
- Map topics = new HashMap<>();
// Add a new partition to topic 1, initially T1 -> 3 partitions
- topics.put(topic1Uuid, new AssignmentTopicMetadata(4));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(3));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 4,
+ createPartitionMetadata(4)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -439,8 +593,8 @@ public void testReassignmentWhenOneConsumerAddedAndOnePartitionAfterInitialAssig
Collections.emptyMap()
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -463,9 +617,26 @@ public void testReassignmentWhenOneConsumerAddedAndOnePartitionAfterInitialAssig
@Test
public void testReassignmentWhenOneConsumerRemovedAfterInitialAssignmentWithTwoConsumersTwoTopics() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(3));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
// Consumer A was removed
@@ -482,8 +653,8 @@ public void testReassignmentWhenOneConsumerRemovedAfterInitialAssignmentWithTwoC
currentAssignmentForB
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -497,10 +668,35 @@ public void testReassignmentWhenOneConsumerRemovedAfterInitialAssignmentWithTwoC
@Test
public void testReassignmentWhenMultipleSubscriptionsRemovedAfterInitialAssignmentWithThreeConsumersTwoTopics() {
- Map topics = new HashMap<>();
- topics.put(topic1Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic2Uuid, new AssignmentTopicMetadata(3));
- topics.put(topic3Uuid, new AssignmentTopicMetadata(2));
+ Map topicMetadata = new HashMap<>();
+ topicMetadata.put(topic1Uuid,
+ new TopicMetadata(
+ topic1Uuid,
+ topic1Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic2Uuid,
+ new TopicMetadata(
+ topic2Uuid,
+ topic2Name,
+ 3,
+ createPartitionMetadata(3)
+ )
+ );
+
+ topicMetadata.put(topic3Uuid,
+ new TopicMetadata(
+ topic3Uuid,
+ topic3Name,
+ 2,
+ createPartitionMetadata(2)
+ )
+ );
+
+ SubscribedTopicMetadata subscribedTopicMetadata = new SubscribedTopicMetadata(topicMetadata);
Map members = new TreeMap<>();
@@ -542,8 +738,8 @@ public void testReassignmentWhenMultipleSubscriptionsRemovedAfterInitialAssignme
currentAssignmentForC
));
- AssignmentSpec assignmentSpec = new AssignmentSpec(members, topics);
- GroupAssignment computedAssignment = assignor.assign(assignmentSpec);
+ AssignmentSpec assignmentSpec = new AssignmentSpec(members);
+ GroupAssignment computedAssignment = assignor.assign(assignmentSpec, subscribedTopicMetadata);
Map>> expectedAssignment = new HashMap<>();
@@ -571,4 +767,13 @@ private void assertAssignment(Map>> expectedAssig
assertEquals(expectedAssignment.get(memberId), computedAssignmentForMember);
}
}
+
+ // When rack awareness is enabled for this assignor, rack information can be updated in this method.
+ private Map> createPartitionMetadata(int numPartitions) {
+ Map> partitionRacks = new HashMap<>(numPartitions);
+ for (int i = 0; i < numPartitions; i++) {
+ partitionRacks.put(i, Collections.emptySet());
+ }
+ return partitionRacks;
+ }
}
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java
index 7306e3a3c5df3..8088af6245fe1 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/ConsumerGroupTest.java
@@ -400,13 +400,13 @@ public void testUpdateSubscriptionMetadata() {
.build();
ConsumerGroupMember member1 = new ConsumerGroupMember.Builder("member1")
- .setSubscribedTopicNames(Arrays.asList("foo"))
+ .setSubscribedTopicNames(Collections.singletonList("foo"))
.build();
ConsumerGroupMember member2 = new ConsumerGroupMember.Builder("member2")
- .setSubscribedTopicNames(Arrays.asList("bar"))
+ .setSubscribedTopicNames(Collections.singletonList("bar"))
.build();
ConsumerGroupMember member3 = new ConsumerGroupMember.Builder("member3")
- .setSubscribedTopicNames(Arrays.asList("zar"))
+ .setSubscribedTopicNames(Collections.singletonList("zar"))
.build();
ConsumerGroup consumerGroup = createConsumerGroup("group-foo");
@@ -417,19 +417,23 @@ public void testUpdateSubscriptionMetadata() {
consumerGroup.computeSubscriptionMetadata(
null,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
// Compute while taking into account member 1.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
null,
member1,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
@@ -439,12 +443,15 @@ public void testUpdateSubscriptionMetadata() {
// It should return foo now.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
null,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
@@ -454,20 +461,26 @@ public void testUpdateSubscriptionMetadata() {
consumerGroup.computeSubscriptionMetadata(
member1,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
// Compute while taking into account member 2.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1)),
- mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ ),
+ mkEntry("bar",
+ new TopicMetadata(barTopicId, "bar", 2, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
null,
member2,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
@@ -477,51 +490,69 @@ public void testUpdateSubscriptionMetadata() {
// It should return foo and bar.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1)),
- mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ ),
+ mkEntry("bar",
+ new TopicMetadata(barTopicId, "bar", 2, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
null,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
// Compute while taking into account removal of member 2.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
member2,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
// Removing member1 results in returning bar.
assertEquals(
mkMap(
- mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2))
+ mkEntry("bar",
+ new TopicMetadata(barTopicId, "bar", 2, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
member1,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
// Compute while taking into account member 3.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1)),
- mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2)),
- mkEntry("zar", new TopicMetadata(zarTopicId, "zar", 3))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ ),
+ mkEntry("bar",
+ new TopicMetadata(barTopicId, "bar", 2, Collections.emptyMap())
+ ),
+ mkEntry("zar",
+ new TopicMetadata(zarTopicId, "zar", 3, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
null,
member3,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
@@ -531,14 +562,21 @@ public void testUpdateSubscriptionMetadata() {
// It should return foo, bar and zar.
assertEquals(
mkMap(
- mkEntry("foo", new TopicMetadata(fooTopicId, "foo", 1)),
- mkEntry("bar", new TopicMetadata(barTopicId, "bar", 2)),
- mkEntry("zar", new TopicMetadata(zarTopicId, "zar", 3))
+ mkEntry("foo",
+ new TopicMetadata(fooTopicId, "foo", 1, Collections.emptyMap())
+ ),
+ mkEntry("bar",
+ new TopicMetadata(barTopicId, "bar", 2, Collections.emptyMap())
+ ),
+ mkEntry("zar",
+ new TopicMetadata(zarTopicId, "zar", 3, Collections.emptyMap())
+ )
),
consumerGroup.computeSubscriptionMetadata(
null,
null,
- image.topics()
+ image.topics(),
+ image.cluster()
)
);
}
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilderTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilderTest.java
index 7733f67b128eb..09195be9bddfd 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilderTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TargetAssignmentBuilderTest.java
@@ -19,7 +19,6 @@
import org.apache.kafka.common.Uuid;
import org.apache.kafka.coordinator.group.assignor.AssignmentMemberSpec;
import org.apache.kafka.coordinator.group.assignor.AssignmentSpec;
-import org.apache.kafka.coordinator.group.assignor.AssignmentTopicMetadata;
import org.apache.kafka.coordinator.group.assignor.GroupAssignment;
import org.apache.kafka.coordinator.group.assignor.MemberAssignment;
import org.apache.kafka.coordinator.group.assignor.PartitionAssignor;
@@ -34,6 +33,8 @@
import java.util.Optional;
import java.util.Set;
+import static org.apache.kafka.common.utils.Utils.mkEntry;
+import static org.apache.kafka.common.utils.Utils.mkMap;
import static org.apache.kafka.coordinator.group.AssignmentTestUtil.mkAssignment;
import static org.apache.kafka.coordinator.group.AssignmentTestUtil.mkTopicAssignment;
import static org.apache.kafka.coordinator.group.RecordHelpers.newTargetAssignmentEpochRecord;
@@ -84,14 +85,15 @@ public void addGroupMember(
public Uuid addTopicMetadata(
String topicName,
- int numPartitions
+ int numPartitions,
+ Map> partitionMetadata
) {
Uuid topicId = Uuid.randomUuid();
subscriptionMetadata.put(topicName, new TopicMetadata(
topicId,
topicName,
- numPartitions
- ));
+ numPartitions,
+ partitionMetadata));
return topicId;
}
@@ -145,13 +147,13 @@ public TargetAssignmentBuilder.TargetAssignmentResult build() {
Map memberSpecs = new HashMap<>();
// All the existing members are prepared.
- members.forEach((memberId, member) -> {
- memberSpecs.put(memberId, createAssignmentMemberSpec(
+ members.forEach((memberId, member) -> memberSpecs.put(memberId,
+ createAssignmentMemberSpec(
member,
targetAssignment.getOrDefault(memberId, Assignment.EMPTY),
subscriptionMetadata
- ));
- });
+ )
+ ));
// All the updated are added and all the deleted
// members are removed.
@@ -168,20 +170,24 @@ public TargetAssignmentBuilder.TargetAssignmentResult build() {
});
// Prepare the expected topic metadata.
- Map topicMetadata = new HashMap<>();
- subscriptionMetadata.forEach((topicName, metadata) -> {
- topicMetadata.put(metadata.id(), new AssignmentTopicMetadata(metadata.numPartitions()));
- });
+ Map topicMetadataMap = new HashMap<>();
+ subscriptionMetadata.forEach((topicName, topicMetadata) ->
+ topicMetadataMap.put(topicMetadata.id(), topicMetadata));
// Prepare the expected assignment spec.
AssignmentSpec assignmentSpec = new AssignmentSpec(
- memberSpecs,
- topicMetadata
+ memberSpecs
+ );
+
+ SubscribedTopicMetadata assignmentTopicMetadata = new SubscribedTopicMetadata(
+ topicMetadataMap
);
// We use `any` here to always return an assignment but use `verify` later on
// to ensure that the input was correct.
- when(assignor.assign(any())).thenReturn(new GroupAssignment(memberAssignments));
+ when(assignor.assign(any(), any()))
+ .thenReturn(new GroupAssignment(memberAssignments));
+
// Create and populate the assignment builder.
TargetAssignmentBuilder builder = new TargetAssignmentBuilder(groupId, groupEpoch, assignor)
@@ -203,7 +209,10 @@ public TargetAssignmentBuilder.TargetAssignmentResult build() {
// Verify that the assignor was called once with the expected
// assignment spec.
- verify(assignor, times(1)).assign(assignmentSpec);
+ verify(assignor, times(1))
+ .assign(
+ assignmentSpec, assignmentTopicMetadata
+ );
return result;
}
@@ -222,8 +231,25 @@ public void testCreateAssignmentMemberSpec() {
Map subscriptionMetadata = new HashMap() {
{
- put("foo", new TopicMetadata(fooTopicId, "foo", 5));
- put("bar", new TopicMetadata(barTopicId, "bar", 5));
+ put("foo", new TopicMetadata(fooTopicId, "foo", 5,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet())
+ )
+ ));
+
+ put("bar", new TopicMetadata(barTopicId, "bar", 5,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet())
+ )
+ ));
}
};
@@ -268,8 +294,27 @@ public void testAssignmentHasNotChanged() {
20
);
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
+ Uuid fooTopicId = context.addTopicMetadata("foo", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
+
+ Uuid barTopicId = context.addTopicMetadata("bar", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"), mkAssignment(
mkTopicAssignment(fooTopicId, 1, 2, 3),
@@ -318,8 +363,27 @@ public void testAssignmentSwapped() {
20
);
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
+ Uuid fooTopicId = context.addTopicMetadata("foo", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
+
+ Uuid barTopicId = context.addTopicMetadata("bar", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"), mkAssignment(
mkTopicAssignment(fooTopicId, 1, 2, 3),
@@ -381,8 +445,27 @@ public void testNewMember() {
20
);
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
+ Uuid fooTopicId = context.addTopicMetadata("foo", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
+
+ Uuid barTopicId = context.addTopicMetadata("bar", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"), mkAssignment(
mkTopicAssignment(fooTopicId, 1, 2, 3),
@@ -459,8 +542,27 @@ public void testUpdateMember() {
20
);
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
+ Uuid fooTopicId = context.addTopicMetadata("foo", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
+
+ Uuid barTopicId = context.addTopicMetadata("bar", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"), mkAssignment(
mkTopicAssignment(fooTopicId, 1, 2, 3),
@@ -546,8 +648,27 @@ public void testPartialAssignmentUpdate() {
20
);
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
+ Uuid fooTopicId = context.addTopicMetadata("foo", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
+
+ Uuid barTopicId = context.addTopicMetadata("bar", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"), mkAssignment(
mkTopicAssignment(fooTopicId, 1, 2),
@@ -624,8 +745,27 @@ public void testDeleteMember() {
20
);
- Uuid fooTopicId = context.addTopicMetadata("foo", 6);
- Uuid barTopicId = context.addTopicMetadata("bar", 6);
+ Uuid fooTopicId = context.addTopicMetadata("foo", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
+
+ Uuid barTopicId = context.addTopicMetadata("bar", 6,
+ mkMap(
+ mkEntry(0, Collections.emptySet()),
+ mkEntry(1, Collections.emptySet()),
+ mkEntry(2, Collections.emptySet()),
+ mkEntry(3, Collections.emptySet()),
+ mkEntry(4, Collections.emptySet()),
+ mkEntry(5, Collections.emptySet())
+ )
+ );
context.addGroupMember("member-1", Arrays.asList("foo", "bar", "zar"), mkAssignment(
mkTopicAssignment(fooTopicId, 1, 2),
diff --git a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TopicMetadataTest.java b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TopicMetadataTest.java
index 07d790ec8ff93..ccbf69821b8e2 100644
--- a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TopicMetadataTest.java
+++ b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/consumer/TopicMetadataTest.java
@@ -20,6 +20,8 @@
import org.apache.kafka.coordinator.group.generated.ConsumerGroupPartitionMetadataValue;
import org.junit.jupiter.api.Test;
+import java.util.Collections;
+
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -28,7 +30,7 @@ public class TopicMetadataTest {
@Test
public void testAttributes() {
Uuid topicId = Uuid.randomUuid();
- TopicMetadata topicMetadata = new TopicMetadata(topicId, "foo", 15);
+ TopicMetadata topicMetadata = new TopicMetadata(topicId, "foo", 15, Collections.emptyMap());
assertEquals(topicId, topicMetadata.id());
assertEquals("foo", topicMetadata.name());
assertEquals(15, topicMetadata.numPartitions());
@@ -36,16 +38,16 @@ public void testAttributes() {
@Test
public void testTopicIdAndNameCannotBeNull() {
- assertThrows(NullPointerException.class, () -> new TopicMetadata(Uuid.randomUuid(), null, 15));
- assertThrows(NullPointerException.class, () -> new TopicMetadata(null, "foo", 15));
+ assertThrows(NullPointerException.class, () -> new TopicMetadata(Uuid.randomUuid(), null, 15, Collections.emptyMap()));
+ assertThrows(NullPointerException.class, () -> new TopicMetadata(null, "foo", 15, Collections.emptyMap()));
}
@Test
public void testEquals() {
Uuid topicId = Uuid.randomUuid();
- TopicMetadata topicMetadata = new TopicMetadata(topicId, "foo", 15);
- assertEquals(new TopicMetadata(topicId, "foo", 15), topicMetadata);
- assertNotEquals(new TopicMetadata(topicId, "foo", 5), topicMetadata);
+ TopicMetadata topicMetadata = new TopicMetadata(topicId, "foo", 15, Collections.emptyMap());
+ assertEquals(new TopicMetadata(topicId, "foo", 15, Collections.emptyMap()), topicMetadata);
+ assertNotEquals(new TopicMetadata(topicId, "foo", 5, Collections.emptyMap()), topicMetadata);
}
@Test
@@ -56,10 +58,11 @@ public void testFromRecord() {
ConsumerGroupPartitionMetadataValue.TopicMetadata record = new ConsumerGroupPartitionMetadataValue.TopicMetadata()
.setTopicId(topicId)
.setTopicName(topicName)
- .setNumPartitions(15);
+ .setNumPartitions(15)
+ .setPartitionMetadata(Collections.emptyList());
assertEquals(
- new TopicMetadata(topicId, topicName, 15),
+ new TopicMetadata(topicId, topicName, 15, Collections.emptyMap()),
TopicMetadata.fromRecord(record)
);
}
diff --git a/reviewers.py b/reviewers.py
old mode 100755
new mode 100644
index 8973a6e862f07..3a00c96989c5a
--- a/reviewers.py
+++ b/reviewers.py
@@ -28,7 +28,7 @@
def prompt_for_user():
while True:
try:
- user_input = input("\nName or email (case insensitive): ")
+ user_input = input("\nName or email (case insensitive): ")
except (KeyboardInterrupt, EOFError):
return None
clean_input = user_input.strip().lower()
@@ -38,7 +38,7 @@ def prompt_for_user():
if __name__ == "__main__":
print("Utility to help generate 'Reviewers' string for Pull Requests. Use Ctrl+D or Ctrl+C to exit")
-
+
stream = os.popen("git log | grep Reviewers")
lines = stream.readlines()
all_reviewers = defaultdict(int)
@@ -46,7 +46,7 @@ def prompt_for_user():
stripped = line.strip().lstrip("Reviewers: ")
reviewers = stripped.split(",")
for reviewer in reviewers:
- all_reviewers[reviewer.strip()] += 1
+ all_reviewers[reviewer.strip()] += 1
parsed_reviewers = []
for item in all_reviewers.items():
@@ -54,7 +54,7 @@ def prompt_for_user():
if m is not None and len(m.groups()) == 2:
if item[1] > 2:
parsed_reviewers.append((m.group("name"), m.group("email"), item[1]))
-
+
selected_reviewers = []
while True:
if selected_reviewers:
@@ -90,5 +90,5 @@ def prompt_for_user():
out += ", ".join([f"{name} <{email}>" for name, email, _ in selected_reviewers])
out += "\n"
print(out)
-
+