diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java index d4e023cc098a8..0ba87c9783d9c 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractStickyAssignor.java @@ -121,6 +121,13 @@ private boolean allSubscriptionsEqual(Set allTopics, 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) { + membersWithOldGeneration.addAll(membersOfCurrentHighestGeneration); + membersOfCurrentHighestGeneration.clear(); + maxGeneration = memberData.generation.get(); + } + 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 @@ -128,13 +135,6 @@ private boolean allSubscriptionsEqual(Set allTopics, ownedPartitions.add(tp); } } - - // If the current member's generation is higher, all the previous owned partitions are invalid - if (memberData.generation.isPresent() && memberData.generation.get() > maxGeneration) { - membersWithOldGeneration.addAll(membersOfCurrentHighestGeneration); - membersOfCurrentHighestGeneration.clear(); - maxGeneration = memberData.generation.get(); - } } }