Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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 @@ -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;
Expand Down Expand Up @@ -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<String> topics, Optional<ConsumerRebalanceListener> 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<ConsumerRebalanceListener> 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<ConsumerRebalanceListener> listener) {
listener.ifPresent(l -> this.listenerContext.set(new ListenerContext(l)));
setSubscriptionType(SubscriptionType.AUTO_PATTERN_RE2J);
this.subscribedRe2JPattern = pattern;
}

public synchronized boolean subscribe(Set<String> topics) {
setSubscriptionType(SubscriptionType.AUTO_TOPICS);
return changeSubscription(topics);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -233,7 +233,7 @@ public void testFirstHeartbeatIncludesRequiredInfoToJoinGroupAndGetAssignments(s
String topic = "topic1";
Set<String> 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);
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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();
Expand All @@ -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();
Expand All @@ -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();
Expand Down Expand Up @@ -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"));

Expand All @@ -201,7 +201,7 @@ public void testNormalSubscription() {
public void testTransientTopics() {
Map<String, Uuid> 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"));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand All @@ -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);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand All @@ -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);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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()));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Loading
Loading