Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,6 @@ private boolean allSubscriptionsEqual(Set<String> allTopics,
Map<String, Subscription> subscriptions,
Map<String, List<TopicPartition>> consumerToOwnedPartitions,
Set<TopicPartition> partitionsWithMultiplePreviousOwners) {
Set<String> membersOfCurrentHighestGeneration = new HashSet<>();
boolean isAllSubscriptionsEqual = true;

Set<String> subscribedTopics = new HashSet<>();
Expand All @@ -106,8 +105,8 @@ private boolean allSubscriptionsEqual(Set<String> allTopics,
Map<TopicPartition, String> allPreviousPartitionsToOwner = new HashMap<>();

for (Map.Entry<String, Subscription> subscriptionEntry : subscriptions.entrySet()) {
String consumer = subscriptionEntry.getKey();
Subscription subscription = subscriptionEntry.getValue();
final String consumer = subscriptionEntry.getKey();
final Subscription subscription = subscriptionEntry.getValue();

// initialize the subscribed topics set if this is the first subscription
if (subscribedTopics.isEmpty()) {
Expand All @@ -118,47 +117,57 @@ private boolean allSubscriptionsEqual(Set<String> allTopics,
}

MemberData memberData = memberData(subscription);
final int memberGeneration = memberData.generation.orElse(DEFAULT_GENERATION);
maxGeneration = Math.max(maxGeneration, memberGeneration);

List<TopicPartition> ownedPartitions = new ArrayList<>();
consumerToOwnedPartitions.put(consumer, ownedPartitions);

// Only consider this consumer's owned partitions as valid if it is a member of the current highest
// generation, or it's generation is not present but we have not seen any known generation so far
if (memberData.generation.isPresent() && memberData.generation.get() >= maxGeneration
|| !memberData.generation.isPresent() && maxGeneration == DEFAULT_GENERATION) {

// If the current member's generation is higher, all the previously owned partitions are invalid
if (memberData.generation.isPresent() && memberData.generation.get() > maxGeneration) {
allPreviousPartitionsToOwner.clear();
partitionsWithMultiplePreviousOwners.clear();
for (String droppedOutConsumer : membersOfCurrentHighestGeneration) {
consumerToOwnedPartitions.get(droppedOutConsumer).clear();
}

membersOfCurrentHighestGeneration.clear();
maxGeneration = memberData.generation.get();
}
// the member has a valid generation, so we can consider its owned partitions if it has the highest
// generation amongst
for (final TopicPartition tp : memberData.partitions) {
if (allTopics.contains(tp.topic())) {
String otherConsumer = allPreviousPartitionsToOwner.put(tp, consumer);
if (otherConsumer == null) {
// this partition is not owned by other consumer in the same generation
ownedPartitions.add(tp);
} else {
final int otherMemberGeneration = subscriptions.get(otherConsumer).generationId().orElse(DEFAULT_GENERATION);

membersOfCurrentHighestGeneration.add(consumer);
for (final TopicPartition tp : memberData.partitions) {
// filter out any topics that no longer exist or aren't part of the current subscription
if (allTopics.contains(tp.topic())) {
String otherConsumer = allPreviousPartitionsToOwner.put(tp, consumer);
if (otherConsumer == null) {
// this partition is not owned by other consumer in the same generation
ownedPartitions.add(tp);
} else {
if (memberGeneration == otherMemberGeneration) {
// if two members of the same generation own the same partition, revoke the partition
log.error("Found multiple consumers {} and {} claiming the same TopicPartition {} in the "
+ "same generation {}, this will be invalidated and removed from their previous assignment.",
consumer, otherConsumer, tp, maxGeneration);
consumerToOwnedPartitions.get(otherConsumer).remove(tp);
+ "same generation {}, this will be invalidated and removed from their previous assignment.",
consumer, otherConsumer, tp, memberGeneration);
partitionsWithMultiplePreviousOwners.add(tp);
consumerToOwnedPartitions.get(otherConsumer).remove(tp);
allPreviousPartitionsToOwner.put(tp, consumer);
} else if (memberGeneration > otherMemberGeneration) {
// move partition from the member with an older generation to the member with the newer generation
ownedPartitions.add(tp);
consumerToOwnedPartitions.get(otherConsumer).remove(tp);
allPreviousPartitionsToOwner.put(tp, consumer);
log.warn("Consumer {} in generation {} and consumer {} in generation {} claiming the same " +
"TopicPartition {} in different generations. The topic partition wil be " +
"assigned to the member with the higher generation {}.",
consumer, memberGeneration,
otherConsumer, otherMemberGeneration,
tp,
memberGeneration);
} else {
// let the other member continue to own the topic partition
log.warn("Consumer {} in generation {} and consumer {} in generation {} claiming the same " +
"TopicPartition {} in different generations. The topic partition wil be " +
"assigned to the member with the higher generation {}.",
consumer, memberGeneration,
otherConsumer, otherMemberGeneration,
tp,
otherMemberGeneration);
}
}
}
}
}

return isAllSubscriptionsEqual;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -942,6 +942,87 @@ public void testPartitionsTransferringOwnershipIncludeThePartitionClaimedByMulti
assertTrue(isFullyBalanced(assignment));
}

public void testEnsurePartitionsAssignedToHighestGeneration() {
Map<String, Integer> partitionsPerTopic = new HashMap<>();
partitionsPerTopic.put(topic, 3);
partitionsPerTopic.put(topic2, 3);
partitionsPerTopic.put(topic3, 3);

int currentGeneration = 10;

// ensure partitions are always assigned to the member with the highest generation
subscriptions.put(consumer1, buildSubscriptionV2Above(topics(topic, topic2, topic3),
partitions(tp(topic, 0), tp(topic2, 0), tp(topic3, 0)), currentGeneration));
subscriptions.put(consumer2, buildSubscriptionV2Above(topics(topic, topic2, topic3),
partitions(tp(topic, 1), tp(topic2, 1), tp(topic3, 1)), currentGeneration - 1));
subscriptions.put(consumer3, buildSubscriptionV2Above(topics(topic, topic2, topic3),
partitions(tp(topic2, 1), tp(topic3, 0), tp(topic3, 2)), currentGeneration - 2));

Map<String, List<TopicPartition>> assignment = assignor.assign(partitionsPerTopic, subscriptions);
assertEquals(new HashSet<>(partitions(tp(topic, 0), tp(topic2, 0), tp(topic3, 0))),
new HashSet<>(assignment.get(consumer1)));
assertEquals(new HashSet<>(partitions(tp(topic, 1), tp(topic2, 1), tp(topic3, 1))),
new HashSet<>(assignment.get(consumer2)));
assertEquals(new HashSet<>(partitions(tp(topic, 2), tp(topic2, 2), tp(topic3, 2))),
new HashSet<>(assignment.get(consumer3)));
assertTrue(assignor.partitionsTransferringOwnership.isEmpty());

verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic);
assertTrue(isFullyBalanced(assignment));
}

public void testNoReassignmentOnCurrentMembers() {
Map<String, Integer> partitionsPerTopic = new HashMap<>();
partitionsPerTopic.put(topic, 3);
partitionsPerTopic.put(topic1, 3);
partitionsPerTopic.put(topic2, 3);
partitionsPerTopic.put(topic3, 3);

int currentGeneration = 10;

subscriptions.put(consumer1, buildSubscriptionV2Above(topics(topic, topic2, topic3, topic1),
partitions(), DEFAULT_GENERATION));
subscriptions.put(consumer2, buildSubscriptionV2Above(topics(topic, topic2, topic3, topic1),
partitions(tp(topic, 0), tp(topic2, 0), tp(topic1, 0)), currentGeneration - 1));
subscriptions.put(consumer3, buildSubscriptionV2Above(topics(topic, topic2, topic3, topic1),
partitions(tp(topic3, 2), tp(topic2, 2), tp(topic1, 1)), currentGeneration - 2));
subscriptions.put(consumer4, buildSubscriptionV2Above(topics(topic, topic2, topic3, topic1),
partitions(tp(topic3, 1), tp(topic, 1), tp(topic, 2)), currentGeneration - 3));

Map<String, List<TopicPartition>> assignment = assignor.assign(partitionsPerTopic, subscriptions);
// ensure assigned partitions don't get reassigned
assertEquals(new HashSet<>(partitions(tp(topic1, 2), tp(topic2, 1), tp(topic3, 0))),
new HashSet<>(assignment.get(consumer1)));
assertTrue(assignor.partitionsTransferringOwnership.isEmpty());

verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic);
assertTrue(isFullyBalanced(assignment));
}

@Test
public void testOwnedPartitionsAreInvalidatedForConsumerWithMultipleGeneration() {
Map<String, Integer> partitionsPerTopic = new HashMap<>();
partitionsPerTopic.put(topic, 3);
partitionsPerTopic.put(topic2, 3);

int currentGeneration = 10;

subscriptions.put(consumer1, buildSubscriptionV2Above(topics(topic, topic2),
partitions(tp(topic, 0), tp(topic2, 1), tp(topic, 1)), currentGeneration));
subscriptions.put(consumer2, buildSubscriptionV2Above(topics(topic, topic2),
partitions(tp(topic, 0), tp(topic2, 1), tp(topic2, 2)), currentGeneration - 2));

Map<String, List<TopicPartition>> assignment = assignor.assign(partitionsPerTopic, subscriptions);
assertEquals(new HashSet<>(partitions(tp(topic, 0), tp(topic2, 1), tp(topic, 1))),
new HashSet<>(assignment.get(consumer1)));
assertEquals(new HashSet<>(partitions(tp(topic, 2), tp(topic2, 2), tp(topic2, 0))),
new HashSet<>(assignment.get(consumer2)));
assertTrue(assignor.partitionsTransferringOwnership.isEmpty());

verifyValidityAndBalance(subscriptions, assignment, partitionsPerTopic);
assertTrue(isFullyBalanced(assignment));
}

private String getTopicName(int i, int maxNum) {
return getCanonicalName("t", i, maxNum);
}
Expand Down