From d809dca4c9f5bd22c0bc3d7e2948037ce0ec1b57 Mon Sep 17 00:00:00 2001 From: Aditya Kousik Date: Wed, 5 Aug 2026 23:09:18 -0700 Subject: [PATCH] KAFKA-20684 [1/N]: Remove test-only SubscriptionState subscribe overloads Drops the three subscribe(X, Optional) overloads and the listener argument at all 52 call sites, all in clients/src/test. --- .../consumer/internals/SubscriptionState.java | 31 ------ .../clients/consumer/KafkaConsumerTest.java | 2 +- .../ConsumerHeartbeatRequestManagerTest.java | 6 +- .../internals/ConsumerMetadataTest.java | 12 +-- .../internals/FetchRequestManagerTest.java | 4 +- .../consumer/internals/FetcherTest.java | 4 +- .../internals/ShareFetchCollectorTest.java | 2 +- .../ShareHeartbeatRequestManagerTest.java | 2 +- .../internals/SubscriptionStateTest.java | 94 +++++++------------ .../ConsumerRebalanceMetricsManagerTest.java | 3 +- 10 files changed, 53 insertions(+), 107 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/SubscriptionState.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/SubscriptionState.java index 32ccf1015cabd..1d3e008e0a9a9 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/SubscriptionState.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/SubscriptionState.java @@ -20,7 +20,6 @@ import org.apache.kafka.clients.Metadata; import org.apache.kafka.clients.NodeApiVersions; import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.apache.kafka.clients.consumer.NoOffsetForPartitionException; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.RebalanceConsumer; @@ -194,36 +193,6 @@ else if (this.subscriptionType != type) throw new IllegalStateException(SUBSCRIPTION_EXCEPTION_MESSAGE); } - /** - * deprecated Visible for testing only. Will be removed in a follow-on cleanup PR. - * Use {@link #subscribe(Set)} and {@link #setRebalanceListener} instead. - */ - public synchronized boolean subscribe(Set topics, Optional listener) { - listener.ifPresent(l -> this.listenerContext.set(new ListenerContext(l))); - setSubscriptionType(SubscriptionType.AUTO_TOPICS); - return changeSubscription(topics); - } - - /** - * deprecated Visible for testing only. Will be removed in a follow-on cleanup PR. - * Use {@link #subscribe(Pattern)} and {@link #setRebalanceListener} instead. - */ - public synchronized void subscribe(Pattern pattern, Optional listener) { - listener.ifPresent(l -> this.listenerContext.set(new ListenerContext(l))); - setSubscriptionType(SubscriptionType.AUTO_PATTERN); - this.subscribedPattern = pattern; - } - - /** - * deprecated Visible for testing only. Will be removed in a follow-on cleanup PR. - * Use {@link #subscribe(SubscriptionPattern)} and {@link #setRebalanceListener} instead. - */ - public synchronized void subscribe(SubscriptionPattern pattern, Optional listener) { - listener.ifPresent(l -> this.listenerContext.set(new ListenerContext(l))); - setSubscriptionType(SubscriptionType.AUTO_PATTERN_RE2J); - this.subscribedRe2JPattern = pattern; - } - public synchronized boolean subscribe(Set topics) { setSubscriptionType(SubscriptionType.AUTO_TOPICS); return changeSubscription(topics); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java index 83a32cc1d85f7..ec45a8d883ec0 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/KafkaConsumerTest.java @@ -303,7 +303,7 @@ public void testAssignedPartitionsMetrics(GroupProtocol groupProtocol) throws In assertEquals(2.0d, getMetric(metrics, "assigned-partitions").metricValue()); subscription.unsubscribe(); - subscription.subscribe(Set.of(topic), Optional.empty()); + subscription.subscribe(Set.of(topic)); subscription.assignFromSubscribed(Set.of(tp0)); assertEquals(1.0d, getMetric(metrics, "assigned-partitions").metricValue()); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java index 065e4d6f56961..fe99af2741edd 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java @@ -233,7 +233,7 @@ public void testFirstHeartbeatIncludesRequiredInfoToJoinGroupAndGetAssignments(s String topic = "topic1"; Set set = Collections.singleton(topic); when(subscriptions.subscription()).thenReturn(set); - subscriptions.subscribe(set, Optional.empty()); + subscriptions.subscribe(set); // Create a ConsumerHeartbeatRequest and verify the payload mockJoiningMemberData(DEFAULT_GROUP_INSTANCE_ID); @@ -612,7 +612,7 @@ public void testHeartbeatState() { // Join the group and subscribe to a topic, but the response has not yet been received String topic = "topic1"; - subscriptions.subscribe(Collections.singleton(topic), Optional.empty()); + subscriptions.subscribe(Collections.singleton(topic)); when(subscriptions.subscription()).thenReturn(Collections.singleton(topic)); mockRejoiningMemberData(); data = heartbeatState.buildRequestData(); @@ -781,7 +781,7 @@ topicId, mkSortedSet(partition) // complete reconciliation createHeartbeatStateAndRequestManager(); when(subscriptions.subscription()).thenReturn(topics); - subscriptions.subscribe(topics, Optional.empty()); + subscriptions.subscribe(topics); mockReconcilingMemberData(testAssignment); // send heartbeat1 to ack assignment tp0 diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerMetadataTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerMetadataTest.java index 49073696d959b..4436ddfe79fcf 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerMetadataTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerMetadataTest.java @@ -80,7 +80,7 @@ public void testPatternSubscriptionIncludeInternalTopics() { } private void testPatternSubscription(boolean includeInternalTopics) { - subscription.subscribe(Pattern.compile("__.*"), Optional.empty()); + subscription.subscribe(Pattern.compile("__.*")); ConsumerMetadata metadata = newConsumerMetadata(includeInternalTopics); MetadataRequest.Builder builder = metadata.newMetadataRequestBuilder(); @@ -104,7 +104,7 @@ private void testPatternSubscription(boolean includeInternalTopics) { @Test public void testSubscriptionToBrokerRegexDoesNotRequestAllTopicsMetadata() { // Subscribe to broker-side regex - subscription.subscribe(new SubscriptionPattern("__.*"), Optional.empty()); + subscription.subscribe(new SubscriptionPattern("__.*")); // Receive assignment from coordinator with topic IDs only Uuid assignedTopicId = Uuid.randomUuid(); @@ -121,7 +121,7 @@ public void testSubscriptionToBrokerRegexDoesNotRequestAllTopicsMetadata() { @Test public void testSubscriptionToBrokerRegexRetainsAssignedTopics() { // Subscribe to broker-side regex - subscription.subscribe(new SubscriptionPattern("__.*"), Optional.empty()); + subscription.subscribe(new SubscriptionPattern("__.*")); // Receive assignment from coordinator with topic IDs only Uuid assignedTopicId = Uuid.randomUuid(); @@ -145,7 +145,7 @@ public void testSubscriptionToBrokerRegexRetainsAssignedTopics() { @Test public void testSubscriptionToBrokerRegexAllowsTransientTopics() { // Subscribe to broker-side regex - subscription.subscribe(new SubscriptionPattern("__.*"), Optional.empty()); + subscription.subscribe(new SubscriptionPattern("__.*")); // Receive assignment from coordinator with topic IDs only Uuid assignedTopicId = Uuid.randomUuid(); @@ -189,7 +189,7 @@ public void testUserAssignment() { @Test public void testNormalSubscription() { - subscription.subscribe(Set.of("foo", "bar", "__consumer_offsets"), Optional.empty()); + subscription.subscribe(Set.of("foo", "bar", "__consumer_offsets")); subscription.groupSubscribe(Set.of("baz", "foo", "bar", "__consumer_offsets")); testBasicSubscription(Set.of("foo", "bar", "baz"), Set.of("__consumer_offsets")); @@ -201,7 +201,7 @@ public void testNormalSubscription() { public void testTransientTopics() { Map topicIds = new HashMap<>(); topicIds.put("foo", Uuid.randomUuid()); - subscription.subscribe(singleton("foo"), Optional.empty()); + subscription.subscribe(singleton("foo")); ConsumerMetadata metadata = newConsumerMetadata(false); metadata.updateWithCurrentRequestVersion(RequestTestUtils.metadataUpdateWithIds(1, singletonMap("foo", 1), topicIds), false, time.milliseconds()); assertEquals(topicIds.get("foo"), metadata.topicIds().get("foo")); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java index f32caebc90907..b0219dd34f2f5 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetchRequestManagerTest.java @@ -1361,7 +1361,7 @@ public void testUnauthorizedTopic() { public void testFetchDuringEagerRebalance() { buildFetcher(); - subscriptions.subscribe(singleton(topicName), Optional.empty()); + subscriptions.subscribe(singleton(topicName)); subscriptions.assignFromSubscribed(singleton(tp0)); subscriptions.seek(tp0, 0); @@ -1385,7 +1385,7 @@ public void testFetchDuringEagerRebalance() { public void testFetchDuringCooperativeRebalance() { buildFetcher(); - subscriptions.subscribe(singleton(topicName), Optional.empty()); + subscriptions.subscribe(singleton(topicName)); subscriptions.assignFromSubscribed(singleton(tp0)); subscriptions.seek(tp0, 0); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java index fe1ecca288e31..7e208900b6a4d 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java @@ -1322,7 +1322,7 @@ public void testUnauthorizedTopic() { public void testFetchDuringEagerRebalance() { buildFetcher(); - subscriptions.subscribe(singleton(topicName), Optional.empty()); + subscriptions.subscribe(singleton(topicName)); subscriptions.assignFromSubscribed(singleton(tp0)); subscriptions.seek(tp0, 0); @@ -1346,7 +1346,7 @@ public void testFetchDuringEagerRebalance() { public void testFetchDuringCooperativeRebalance() { buildFetcher(); - subscriptions.subscribe(singleton(topicName), Optional.empty()); + subscriptions.subscribe(singleton(topicName)); subscriptions.assignFromSubscribed(singleton(tp0)); subscriptions.seek(tp0, 0); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java index afcafe922f551..48ef9098226ef 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java @@ -335,7 +335,7 @@ private void buildDependencies() { } private void subscribeAndAssign(TopicIdPartition tp) { - subscriptions.subscribe(Set.of(tp.topic()), Optional.empty()); + subscriptions.subscribe(Set.of(tp.topic())); subscriptions.assignFromSubscribed(Set.of(tp.topicPartition())); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java index b31a834e4e159..1eebd52d550f5 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java @@ -376,7 +376,7 @@ public void testHeartbeatState() { // Join the group and subscribe to a topic, but the response has not yet been received String topic = "topic1"; - subscriptions.subscribe(Set.of(topic), Optional.empty()); + subscriptions.subscribe(Set.of(topic)); when(subscriptions.subscription()).thenReturn(Set.of(topic)); mockRejoiningMemberData(); data = heartbeatState.buildRequestData(); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/SubscriptionStateTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/SubscriptionStateTest.java index 27846116262bd..02f42377319f5 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/SubscriptionStateTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/SubscriptionStateTest.java @@ -66,7 +66,6 @@ public class SubscriptionStateTest { private final TopicPartition tp0 = new TopicPartition(topic, 0); private final TopicPartition tp1 = new TopicPartition(topic, 1); private final TopicPartition t1p0 = new TopicPartition(topic1, 0); - private final MockRebalanceListener rebalanceListener = new MockRebalanceListener(); private final Metadata.LeaderAndEpoch leaderAndEpoch = Metadata.LeaderAndEpoch.noLeaderOrEpoch(); private final Collection partitions = List.of(tp0, tp1); @@ -107,7 +106,7 @@ public void partitionAssignmentChangeOnTopicSubscription() { assertTrue(state.assignedPartitions().isEmpty()); assertEquals(0, state.numAssignedPartitions()); - state.subscribe(Set.of(topic1), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic1)); // assigned partitions should remain unchanged assertTrue(state.assignedPartitions().isEmpty()); assertEquals(0, state.numAssignedPartitions()); @@ -118,7 +117,7 @@ public void partitionAssignmentChangeOnTopicSubscription() { assertEquals(Set.of(t1p0), state.assignedPartitions()); assertEquals(1, state.numAssignedPartitions()); - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); // assigned partitions should remain unchanged assertEquals(Set.of(t1p0), state.assignedPartitions()); assertEquals(1, state.numAssignedPartitions()); @@ -137,7 +136,7 @@ public void testIsFetchableOnManualAssignment() { @Test public void testIsFetchableOnAutoAssignment() { - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); state.assignFromSubscribed(Set.of(tp0, tp1)); assertAssignedPartitionIsFetchable(); } @@ -159,7 +158,7 @@ private void assertAssignedPartitionIsFetchable() { @Test public void testIsFetchableConsidersExplicitTopicSubscription() { - state.subscribe(Set.of(topic1), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic1)); state.assignFromSubscribed(Set.of(t1p0)); state.seek(t1p0, 1); @@ -167,7 +166,7 @@ public void testIsFetchableConsidersExplicitTopicSubscription() { assertTrue(state.isFetchable(t1p0)); // Change subscription. Assigned partitions should remain unchanged but not fetchable. - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); assertEquals(Set.of(t1p0), state.assignedPartitions()); assertFalse(state.isFetchable(t1p0), "Assigned partitions not in the subscription should not be fetchable"); @@ -179,7 +178,7 @@ public void testIsFetchableConsidersExplicitTopicSubscription() { @Test public void testGroupSubscribe() { - state.subscribe(Set.of(topic1), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic1)); assertEquals(Set.of(topic1), state.metadataTopics()); assertFalse(state.groupSubscribe(Set.of(topic1))); @@ -192,7 +191,7 @@ public void testGroupSubscribe() { assertFalse(state.groupSubscribe(Set.of(topic1))); assertEquals(Set.of(topic1), state.metadataTopics()); - state.subscribe(Set.of("anotherTopic"), Optional.of(rebalanceListener)); + state.subscribe(Set.of("anotherTopic")); assertEquals(Set.of(topic1, "anotherTopic"), state.metadataTopics()); assertFalse(state.groupSubscribe(Set.of("anotherTopic"))); @@ -201,7 +200,7 @@ public void testGroupSubscribe() { @Test public void partitionAssignmentChangeOnPatternSubscription() { - state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener)); + state.subscribe(Pattern.compile(".*")); // assigned partitions should remain unchanged assertTrue(state.assignedPartitions().isEmpty()); assertEquals(0, state.numAssignedPartitions()); @@ -227,7 +226,7 @@ public void partitionAssignmentChangeOnPatternSubscription() { assertEquals(1, state.numAssignedPartitions()); assertEquals(Set.of(topic), state.subscription()); - state.subscribe(Pattern.compile(".*t"), Optional.of(rebalanceListener)); + state.subscribe(Pattern.compile(".*t")); // assigned partitions should remain unchanged assertEquals(Set.of(t1p0), state.assignedPartitions()); assertEquals(1, state.numAssignedPartitions()); @@ -264,7 +263,7 @@ public void verifyAssignmentId() { assertEquals(Set.of(), state.assignedPartitions()); Set autoAssignment = Set.of(t1p0); - state.subscribe(Set.of(topic1), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic1)); assertTrue(state.checkAssignmentMatchedSubscription(autoAssignment)); state.assignFromSubscribed(autoAssignment); assertEquals(3, state.assignmentId()); @@ -289,7 +288,7 @@ public void partitionReset() { @Test public void topicSubscription() { - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); assertEquals(1, state.subscription().size()); assertTrue(state.assignedPartitions().isEmpty()); assertEquals(0, state.numAssignedPartitions()); @@ -342,7 +341,7 @@ public void testMarkingPendingRevocationPreventsInitializingPosition() { @Test public void testAssignedPartitionsAwaitingCallbackKeepPositionDefinedInCallback() { // New partition assigned. Should not be fetchable or initializing positions. - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); state.assignFromSubscribedAwaitingCallback(Set.of(tp0), Set.of(tp0)); assertAssignmentAppliedAwaitingCallback(tp0); assertEquals(Set.of(tp0.topic()), state.subscription()); @@ -362,7 +361,7 @@ public void testAssignedPartitionsAwaitingCallbackKeepPositionDefinedInCallback( @Test public void testAssignedPartitionsAwaitingCallbackInitializePositionsWhenCallbackCompletes() { // New partition assigned. Should not be fetchable or initializing positions. - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); state.assignFromSubscribedAwaitingCallback(Set.of(tp0), Set.of(tp0)); assertAssignmentAppliedAwaitingCallback(tp0); assertEquals(Set.of(tp0.topic()), state.subscription()); @@ -380,7 +379,7 @@ public void testAssignedPartitionsAwaitingCallbackInitializePositionsWhenCallbac @Test public void testAssignedPartitionsAwaitingCallbackDoesNotAffectPreviouslyOwnedPartitions() { // First partition assigned and callback completes. - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); state.assignFromSubscribedAwaitingCallback(Set.of(tp0), Set.of(tp0)); assertAssignmentAppliedAwaitingCallback(tp0); assertEquals(Set.of(tp0.topic()), state.subscription()); @@ -414,7 +413,7 @@ private void assertAssignmentAppliedAwaitingCallback(TopicPartition topicPartiti @Test public void invalidPositionUpdate() { - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); assertTrue(state.checkAssignmentMatchedSubscription(Set.of(tp0))); state.assignFromSubscribed(Set.of(tp0)); @@ -424,13 +423,13 @@ public void invalidPositionUpdate() { @Test public void cantAssignPartitionForUnsubscribedTopics() { - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); assertFalse(state.checkAssignmentMatchedSubscription(List.of(t1p0))); } @Test public void cantAssignPartitionForUnmatchedPattern() { - state.subscribe(Pattern.compile(".*t"), Optional.of(rebalanceListener)); + state.subscribe(Pattern.compile(".*t")); state.subscribeFromPattern(Set.of(topic)); assertFalse(state.checkAssignmentMatchedSubscription(List.of(t1p0))); } @@ -443,31 +442,31 @@ public void cantChangePositionForNonAssignedPartition() { @Test public void cantSubscribeTopicAndPattern() { - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); - assertThrows(IllegalStateException.class, () -> state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener))); + state.subscribe(Set.of(topic)); + assertThrows(IllegalStateException.class, () -> state.subscribe(Pattern.compile(".*"))); } @Test public void cantSubscribePartitionAndPattern() { state.assignFromUser(Set.of(tp0)); - assertThrows(IllegalStateException.class, () -> state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener))); + assertThrows(IllegalStateException.class, () -> state.subscribe(Pattern.compile(".*"))); } @Test public void cantSubscribePatternAndTopic() { - state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener)); - assertThrows(IllegalStateException.class, () -> state.subscribe(Set.of(topic), Optional.of(rebalanceListener))); + state.subscribe(Pattern.compile(".*")); + assertThrows(IllegalStateException.class, () -> state.subscribe(Set.of(topic))); } @Test public void cantSubscribePatternAndPartition() { - state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener)); + state.subscribe(Pattern.compile(".*")); assertThrows(IllegalStateException.class, () -> state.assignFromUser(Set.of(tp0))); } @Test public void patternSubscription() { - state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener)); + state.subscribe(Pattern.compile(".*")); state.subscribeFromPattern(Set.of(topic, topic1)); assertEquals(2, state.subscription().size(), "Expected subscribed topics count is incorrect"); } @@ -475,7 +474,7 @@ public void patternSubscription() { @Test public void testSubscribeToRe2JPattern() { String pattern = "t.*"; - state.subscribe(new SubscriptionPattern(pattern), Optional.of(rebalanceListener)); + state.subscribe(new SubscriptionPattern(pattern)); assertTrue(state.toString().contains("type=AUTO_PATTERN_RE2J")); assertTrue(state.toString().contains("subscribedPattern=" + pattern)); assertTrue(state.assignedTopicIds().isEmpty()); @@ -487,7 +486,7 @@ public void testIsAssignedFromRe2j() { Uuid assignedUuid = Uuid.randomUuid(); assertFalse(state.isAssignedFromRe2j(assignedUuid)); - state.subscribe(new SubscriptionPattern("foo.*"), Optional.empty()); + state.subscribe(new SubscriptionPattern("foo.*")); assertTrue(state.hasRe2JPatternSubscription()); assertFalse(state.isAssignedFromRe2j(assignedUuid)); @@ -502,7 +501,7 @@ public void testIsAssignedFromRe2j() { @Test public void testAssignedPartitionsWithTopicIdsForRe2Pattern() { - state.subscribe(new SubscriptionPattern("t.*"), Optional.of(rebalanceListener)); + state.subscribe(new SubscriptionPattern("t.*")); assertTrue(state.assignedTopicIds().isEmpty()); TopicIdPartitionSet reconciledAssignmentFromRegex = new TopicIdPartitionSet(); @@ -523,7 +522,7 @@ public void testAssignedPartitionsWithTopicIdsForRe2Pattern() { @Test public void testAssignedTopicIdsPreservedWhenReconciliationCompletes() { - state.subscribe(new SubscriptionPattern("t.*"), Optional.of(rebalanceListener)); + state.subscribe(new SubscriptionPattern("t.*")); assertTrue(state.assignedTopicIds().isEmpty()); // First assignment received from coordinator @@ -551,20 +550,19 @@ public void testAssignedTopicIdsPreservedWhenReconciliationCompletes() { @Test public void testMixedPatternSubscriptionNotAllowed() { - state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener)); - assertThrows(IllegalStateException.class, () -> state.subscribe(new SubscriptionPattern("t.*"), - Optional.of(rebalanceListener))); + state.subscribe(Pattern.compile(".*")); + assertThrows(IllegalStateException.class, () -> state.subscribe(new SubscriptionPattern("t.*"))); state.unsubscribe(); - state.subscribe(new SubscriptionPattern("t.*"), Optional.of(rebalanceListener)); - assertThrows(IllegalStateException.class, () -> state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener))); + state.subscribe(new SubscriptionPattern("t.*")); + assertThrows(IllegalStateException.class, () -> state.subscribe(Pattern.compile(".*"))); } @Test public void testSubscriptionPattern() { SubscriptionPattern pattern = new SubscriptionPattern("t.*"); - state.subscribe(pattern, Optional.of(rebalanceListener)); + state.subscribe(pattern); assertTrue(state.hasRe2JPatternSubscription()); assertEquals(pattern, state.subscriptionPattern()); assertTrue(state.hasAutoAssignedPartitions()); @@ -579,13 +577,13 @@ public void testSubscriptionPattern() { public void unsubscribeUserAssignment() { state.assignFromUser(Set.of(tp0, tp1)); state.unsubscribe(); - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); assertEquals(Set.of(topic), state.subscription()); } @Test public void unsubscribeUserSubscribe() { - state.subscribe(Set.of(topic), Optional.of(rebalanceListener)); + state.subscribe(Set.of(topic)); state.unsubscribe(); state.assignFromUser(Set.of(tp0)); assertEquals(Set.of(tp0), state.assignedPartitions()); @@ -594,7 +592,7 @@ public void unsubscribeUserSubscribe() { @Test public void unsubscription() { - state.subscribe(Pattern.compile(".*"), Optional.of(rebalanceListener)); + state.subscribe(Pattern.compile(".*")); state.subscribeFromPattern(Set.of(topic, topic1)); assertTrue(state.checkAssignmentMatchedSubscription(Set.of(tp1))); state.assignFromSubscribed(Set.of(tp1)); @@ -1208,26 +1206,6 @@ public void onPartitionsRevoked(Collection partitions) { assertEquals(List.of("assigned-1arg", "revoked-1arg", "revoked-1arg"), calls); } - private static class MockRebalanceListener implements ConsumerRebalanceListener { - Collection revoked; - public Collection assigned; - int revokedCount = 0; - int assignedCount = 0; - - @Override - public void onPartitionsAssigned(Collection partitions) { - this.assigned = partitions; - assignedCount++; - } - - @Override - public void onPartitionsRevoked(Collection partitions) { - this.revoked = partitions; - revokedCount++; - } - - } - @Test public void resetOffsetNoValidation() { // Check that offset reset works when we can't validate offsets (older brokers) diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/ConsumerRebalanceMetricsManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/ConsumerRebalanceMetricsManagerTest.java index 639ba823f3566..ed882ddd3cae4 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/ConsumerRebalanceMetricsManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/metrics/ConsumerRebalanceMetricsManagerTest.java @@ -29,7 +29,6 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import java.util.Optional; import java.util.Set; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -88,7 +87,7 @@ public void testAssignedPartitionCountMetric() { assertEquals(0.0d, metrics.metric(metricsManager.assignedPartitionsCount).metricValue()); // Check for automatically assigned partitions - subscriptionState.subscribe(Set.of("topic"), Optional.empty()); + subscriptionState.subscribe(Set.of("topic")); subscriptionState.assignFromSubscribed(Set.of(new TopicPartition("topic", 0))); assertEquals(1.0d, metrics.metric(metricsManager.assignedPartitionsCount).metricValue()); }