From 05ddb3be068c2648ed72ae2b04c269c87a901dfb Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 23 Jul 2020 16:30:17 -0700 Subject: [PATCH 1/2] first fix --- .../consumer/internals/ConsumerCoordinator.java | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java index 9a932f9ac15d9..12657b85bec11 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java @@ -770,12 +770,18 @@ public boolean rejoinNeededOrPending() { // we need to rejoin if we performed the assignment and metadata has changed; // also for those owned-but-no-longer-existed partitions we should drop them as lost if (assignmentSnapshot != null && !assignmentSnapshot.matches(metadataSnapshot)) { + log.info("Requesting to re-join the group and trigger rebalance since the assignment metadata has changed from {} to {}", + assignmentSnapshot, metadataSnapshot); + requestRejoin(); return true; } // we need to join if our subscription has changed since the last join if (joinedSubscription != null && !joinedSubscription.equals(subscriptions.subscription())) { + log.info("Requesting to re-join the group and trigger rebalance since the subscription has changed from {} to {}", + joinedSubscription, subscriptions.subscription()); + requestRejoin(); return true; } @@ -1436,6 +1442,10 @@ boolean matches(MetadataSnapshot other) { return version == other.version || partitionsPerTopic.equals(other.partitionsPerTopic); } + @Override + public String toString() { + return "(version" + version + ": " + partitionsPerTopic + ")"; + } } private static class OffsetCommitCompletion { From 5edcb4aca811a640f83027fdf805911f9e9702c2 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Sun, 26 Jul 2020 16:40:48 -0700 Subject: [PATCH 2/2] github comments --- .../kafka/clients/consumer/internals/ConsumerCoordinator.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java index 12657b85bec11..fd58a60740a56 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java @@ -780,7 +780,7 @@ public boolean rejoinNeededOrPending() { // we need to join if our subscription has changed since the last join if (joinedSubscription != null && !joinedSubscription.equals(subscriptions.subscription())) { log.info("Requesting to re-join the group and trigger rebalance since the subscription has changed from {} to {}", - joinedSubscription, subscriptions.subscription()); + joinedSubscription, subscriptions.subscription()); requestRejoin(); return true;