Skip to content
Merged
Show file tree
Hide file tree
Changes from 21 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 @@ -127,7 +127,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 @@ -137,8 +136,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 @@ -149,47 +148,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) {
Comment thread
philipnee marked this conversation as resolved.
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) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Could we put a comment in this branch like we did for the others?

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
consumerToOwnedPartitions.get(consumer).add(tp);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: We could use ownedPartitions.add(tp).

consumerToOwnedPartitions.get(otherConsumer).remove(tp);
allPreviousPartitionsToOwner.put(tp, consumer);
// if memberGeneration > otherMemberGeneration, the other member continue owns the generation

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Should we just remove this one? It seems that the comment earlier explain it all.

log.warn("{} in generation {} and {} in generation {} claiming the same TopicPartition {} in " +

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: I would say Consumer {} in generation and consumer {} in .... This applies to the other comment as well.

"different generations. The topic partition wil be assigned to the member with " +
"the higher generation {}.",
consumer, memberGeneration,
otherConsumer, otherMemberGeneration,
tp,
memberGeneration);
} else {
// if memberGeneration < otherMemberGeneration, the other member continue owns the generation

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Should we remove if memberGeneration < otherMemberGeneration as it is implicit?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks, there's also a typo: continue to own the "topic partition"

log.warn("{} in generation {} and {} 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 @@ -1038,6 +1038,99 @@ public void testPartitionsTransferringOwnershipIncludeThePartitionClaimedByMulti
assertTrue(isFullyBalanced(assignment));
}

@ParameterizedTest(name = TEST_NAME_WITH_RACK_CONFIG)
@EnumSource(RackConfig.class)
public void testEnsurePartitionsAssignedToHighestGeneration(RackConfig rackConfig) {
initializeRacks(rackConfig);
Map<String, List<PartitionInfo>> partitionsPerTopic = new HashMap<>();
partitionsPerTopic.put(topic, partitionInfos(topic, 3));
partitionsPerTopic.put(topic2, partitionInfos(topic2, 3));
partitionsPerTopic.put(topic3, partitionInfos(topic3, 3));

int currentGeneration = 10;

// ensure partitions are always assigned to the member with the highest generation
// topic, 1 -> [consumer2], consumer3
// topic2, 1 -> [consumer2], consumer3
// topic, 0 -> [consumer1], consumer3
// topic3, 2 -> consumer3

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't understand what you are trying to express here. It also seems to me that the mapping is incorrect. Should we just remove this?

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Could we put tp(topic2, 1) before topic3?


Map<String, List<TopicPartition>> assignment = assignor.assignPartitions(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));
}

@ParameterizedTest(name = TEST_NAME_WITH_RACK_CONFIG)
@EnumSource(RackConfig.class)
public void testNoReassignmentOnCurrentMembers(RackConfig rackConfig) {
initializeRacks(rackConfig);
Map<String, List<PartitionInfo>> partitionsPerTopic = new HashMap<>();
partitionsPerTopic.put(topic, partitionInfos(topic, 3));
partitionsPerTopic.put(topic2, partitionInfos(topic2, 3));
partitionsPerTopic.put(topic3, partitionInfos(topic3, 3));
partitionsPerTopic.put(topic1, partitionInfos(topic1, 3));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Could we put topic1 after topic?


int currentGeneration = 10;

subscriptions.put(consumer1, buildSubscriptionV2Above(topics(topic, topic2, topic3, topic1), partitions(),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Could we put partitions(), on the next line to follow the pattern used by the others? It is easier to read and compare them.

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

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: We should put the expected one first.

assertTrue(assignor.partitionsTransferringOwnership.isEmpty());

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

@ParameterizedTest(name = TEST_NAME_WITH_RACK_CONFIG)
@EnumSource(RackConfig.class)
public void testOwnedPartitionsAreInvalidatedForConsumerWithMultipleGeneration(RackConfig rackConfig) {
initializeRacks(rackConfig);
Map<String, List<PartitionInfo>> partitionsPerTopic = new HashMap<>();
partitionsPerTopic.put(topic, partitionInfos(topic, 3));
partitionsPerTopic.put(topic2, partitionInfos(topic2, 3));

int currentGeneration = 10;

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

Map<String, List<TopicPartition>> assignment = assignor.assignPartitions(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));
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we also add a test for the case where a partition is claimed by two consumers with the same generation?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

good idea, I didn't check because i thought there could be an existing test case covering it but let me check

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ok testPartitionsTransferringOwnershipIncludeThePartitionClaimedByMultipleConsumersInSameGeneration might cover it. But I'll double check if that's what we want.

@Test
public void testRackAwareAssignmentWithUniformSubscription() {
Map<String, Integer> topics = mkMap(mkEntry("t1", 6), mkEntry("t2", 7), mkEntry("t3", 2));
Expand Down