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..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 @@ -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 {