Skip to content
Open
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 @@ -17,10 +17,11 @@
package org.apache.kafka.clients;

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.consumer.OffsetCommitCallback;
import org.apache.kafka.clients.consumer.RebalanceConsumer;
import org.apache.kafka.clients.consumer.RebalanceListener;
import org.apache.kafka.clients.consumer.RetriableCommitFailedException;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
Expand Down Expand Up @@ -434,7 +435,8 @@ public static void testCoordinatorFailover(
) throws InterruptedException {
var listener = new TestConsumerReassignmentListener();
try (Consumer<byte[], byte[]> consumer = cluster.consumer(consumerConfig)) {
consumer.subscribe(List.of(TOPIC), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(TOPIC));
// the initial subscription should cause a callback execution
awaitRebalance(consumer, listener);
assertEquals(1, listener.callsToAssigned);
Expand Down Expand Up @@ -531,17 +533,17 @@ public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets, Exception
}
}

public static class TestConsumerReassignmentListener implements ConsumerRebalanceListener {
public static class TestConsumerReassignmentListener implements RebalanceListener {
public int callsToAssigned = 0;
public int callsToRevoked = 0;

@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer consumer) {
callsToAssigned += 1;
}

@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer consumer) {
callsToRevoked += 1;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ public class ConsumerAssignmentPoller extends ShutdownableThread {
private final Set<TopicPartition> partitionAssignment = Collections.synchronizedSet(new HashSet<>());
private volatile boolean subscriptionChanged = false;
private List<String> topicsSubscription;
private final ConsumerRebalanceListener rebalanceListener;
private final RebalanceListener rebalanceListener;

public ConsumerAssignmentPoller(Consumer<byte[], byte[]> consumer, List<String> topicsToSubscribe) {
this(consumer, topicsToSubscribe, Set.of(), null);
Expand All @@ -49,30 +49,31 @@ public ConsumerAssignmentPoller(Consumer<byte[], byte[]> consumer, Set<TopicPart
public ConsumerAssignmentPoller(Consumer<byte[], byte[]> consumer,
List<String> topicsToSubscribe,
Set<TopicPartition> partitionsToAssign,
ConsumerRebalanceListener userRebalanceListener) {
RebalanceListener userRebalanceListener) {
super("daemon-consumer-assignment", false);
this.consumer = consumer;
this.partitionsToAssign = partitionsToAssign;
this.topicsSubscription = topicsToSubscribe;

this.rebalanceListener = new ConsumerRebalanceListener() {
this.rebalanceListener = new RebalanceListener() {
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
partitionAssignment.addAll(partitions);
if (userRebalanceListener != null)
userRebalanceListener.onPartitionsAssigned(partitions);
userRebalanceListener.onPartitionsAssigned(partitions, rebalanceConsumer);
}

@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
partitionAssignment.removeAll(partitions);
if (userRebalanceListener != null)
userRebalanceListener.onPartitionsRevoked(partitions);
userRebalanceListener.onPartitionsRevoked(partitions, rebalanceConsumer);
}
};

if (partitionsToAssign.isEmpty()) {
consumer.subscribe(topicsToSubscribe, rebalanceListener);
consumer.setRebalanceListener(rebalanceListener);
consumer.subscribe(topicsToSubscribe);
} else {
consumer.assign(List.copyOf(partitionsToAssign));
}
Expand Down Expand Up @@ -107,7 +108,8 @@ public boolean initiateShutdown() {
@Override
public void doWork() {
if (subscriptionChanged) {
consumer.subscribe(topicsSubscription, rebalanceListener);
consumer.setRebalanceListener(rebalanceListener);
consumer.subscribe(topicsSubscription);
subscriptionChanged = false;
}
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -724,17 +724,18 @@ private void checkClosedState(String groupId, int committedRecords) throws Inter

Semaphore assignSemaphore = new Semaphore(0);
try (Consumer<byte[], byte[]> consumer = clusterInstance.consumer(Map.of(ConsumerConfig.GROUP_ID_CONFIG, groupId))) {
consumer.subscribe(List.of(topic), new ConsumerRebalanceListener() {
consumer.setRebalanceListener(new RebalanceListener() {
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
assignSemaphore.release();
}

@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
// Do nothing
}
});
consumer.subscribe(List.of(topic));

TestUtils.waitForCondition(() -> {
consumer.poll(Duration.ofMillis(100));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -116,18 +116,19 @@ private static void testFetchPartitionsAfterFailedListener(ClusterInstance clust

try (var consumer = clusterInstance.consumer(Map.of(
ConsumerConfig.GROUP_PROTOCOL_CONFIG, groupProtocol.name()))) {
consumer.subscribe(List.of(topic), new ConsumerRebalanceListener() {
consumer.setRebalanceListener(new RebalanceListener() {
private int count = 0;
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
}

@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
count++;
if (count == 1) throw new IllegalArgumentException("temporary error");
}
});
consumer.subscribe(List.of(topic));

TestUtils.waitForCondition(() -> consumer.poll(Duration.ofSeconds(1)).count() == 1,
5000,
Expand Down Expand Up @@ -164,16 +165,17 @@ private static void testFetchPartitionsWithAlwaysFailedListener(ClusterInstance

try (var consumer = clusterInstance.consumer(Map.of(
ConsumerConfig.GROUP_PROTOCOL_CONFIG, groupProtocol.name()))) {
consumer.subscribe(List.of(topic), new ConsumerRebalanceListener() {
consumer.setRebalanceListener(new RebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
}

@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
throw new IllegalArgumentException("always failed");
}
});
consumer.subscribe(List.of(topic));

long startTimeMillis = System.currentTimeMillis();
long currentTimeMillis = System.currentTimeMillis();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -334,14 +334,14 @@ private void testAutoCommitIntercept(GroupProtocol groupProtocol) throws Interru
producer.send(new ProducerRecord<>(tp.topic(), tp.partition(), ("key " + i).getBytes(), ("value " + i).getBytes()));
}

var rebalanceListener = new ConsumerRebalanceListener() {
var rebalanceListener = new RebalanceListener() {
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
// keep partitions paused in this test so that we can verify the commits based on specific seeks
consumer.pause(partitions);
}
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
// No-op
}
};
Expand Down Expand Up @@ -448,26 +448,28 @@ private void testAutoCommitOnRebalance(GroupProtocol groupProtocol) throws Inter
try (var consumer = createConsumer(groupProtocol, true)) {
sendRecords(cluster, tp, 1000);

var rebalanceListener = new ConsumerRebalanceListener() {
var rebalanceListener = new RebalanceListener() {
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
// keep partitions paused in this test so that we can verify the commits based on specific seeks
consumer.pause(partitions);
}

@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {

}
};

consumer.subscribe(List.of(topic), rebalanceListener);
consumer.setRebalanceListener(rebalanceListener);
consumer.subscribe(List.of(topic));
awaitAssignment(consumer, Set.of(tp, tp1));

consumer.seek(tp, 300);
consumer.seek(tp1, 500);
// change subscription to trigger rebalance
consumer.subscribe(List.of(topic, topic2), rebalanceListener);
consumer.setRebalanceListener(rebalanceListener);
consumer.subscribe(List.of(topic, topic2));

var newAssignment = Set.of(tp, tp1, new TopicPartition(topic2, 0), new TopicPartition(topic2, 1));
awaitAssignment(consumer, newAssignment);
Expand Down Expand Up @@ -719,9 +721,10 @@ private void changeConsumerSubscriptionAndValidateAssignment(
Consumer<byte[], byte[]> consumer,
List<String> topicsToSubscribe,
Set<TopicPartition> expectedAssignment,
ConsumerRebalanceListener rebalanceListener
RebalanceListener rebalanceListener
) throws InterruptedException {
consumer.subscribe(topicsToSubscribe, rebalanceListener);
consumer.setRebalanceListener(rebalanceListener);
consumer.subscribe(topicsToSubscribe);
awaitAssignment(consumer, expectedAssignment);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,8 @@ public void testAsyncConsumerMaxPollIntervalMs() throws InterruptedException {
private void testMaxPollIntervalMs(Map<String, Object> config) throws InterruptedException {
try (Consumer<byte[], byte[]> consumer = cluster.consumer(config)) {
var listener = new TestConsumerReassignmentListener();
consumer.subscribe(List.of(topic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(topic));

// rebalance to get the initial assignment
awaitRebalance(consumer, listener);
Expand Down Expand Up @@ -203,12 +204,12 @@ private void testMaxPollIntervalMsDelayInRevocation(Map<String, Object> config)
try (Consumer<byte[], byte[]> consumer = cluster.consumer(config)) {
var listener = new TestConsumerReassignmentListener() {
@Override
public void onPartitionsLost(Collection<TopicPartition> partitions) {
public void onPartitionsLost(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
// no op
}

@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
if (!partitions.isEmpty() && partitions.contains(tp)) {
// on the second rebalance (after we have joined the group initially), sleep longer
// than session timeout and then try a commit. We should still be in the group,
Expand All @@ -219,17 +220,19 @@ public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
consumer.commitSync(offsets);
commitCompleted.set(true);
}
super.onPartitionsRevoked(partitions);
super.onPartitionsRevoked(partitions, rebalanceConsumer);
}
};
consumer.subscribe(List.of(topic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(topic));

// Consume records to ensure the rebalance completed and positions are initialized
// (position then used in callback, triggered on next rebalance)
awaitNonEmptyRecords(consumer, tp, 100);

// force a rebalance to trigger an invocation of the revocation callback while in the group
consumer.subscribe(List.of(otherTopic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(otherTopic));
// Consume records to ensure positions for otherTopic are initialized
// (position then used in callback, triggered on close)
awaitNonEmptyRecords(consumer, tpOther, 100);
Expand Down Expand Up @@ -263,13 +266,14 @@ private void testMaxPollIntervalMsDelayInAssignment(Map<String, Object> config)
try (Consumer<byte[], byte[]> consumer = cluster.consumer(config)) {
var listener = new TestConsumerReassignmentListener() {
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
// sleep longer than the session timeout, we should still be in the group after invocation
Utils.sleep(1500);
super.onPartitionsAssigned(partitions);
super.onPartitionsAssigned(partitions, rebalanceConsumer);
}
};
consumer.subscribe(List.of(topic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(topic));
// rebalance to get the initial assignment
awaitRebalance(consumer, listener);
// We should still be in the group after this invocation
Expand Down Expand Up @@ -297,7 +301,8 @@ public void testAsyncConsumerMaxPollIntervalMsShorterThanPollTimeout() throws In
private void testMaxPollIntervalMsShorterThanPollTimeout(Map<String, Object> config) throws InterruptedException {
try (Consumer<byte[], byte[]> consumer = cluster.consumer(config)) {
var listener = new TestConsumerReassignmentListener();
consumer.subscribe(List.of(topic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(topic));

// rebalance to get the initial assignment
awaitRebalance(consumer, listener);
Expand Down Expand Up @@ -551,23 +556,25 @@ public void testConsumerRecoveryOnPollAfterDelayedRebalance(GroupProtocol groupP

var listener = new TestConsumerReassignmentListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(Collection<TopicPartition> partitions, RebalanceConsumer rebalanceConsumer) {
if (!partitions.isEmpty() && partitions.contains(tp)) {
// on the second rebalance (after we have joined the group initially), sleep longer
// than rebalance timeout to get fenced.
Utils.sleep(rebalanceTimeout + 500);
rebalanceTimeoutExceeded.set(true);
}
super.onPartitionsRevoked(partitions);
super.onPartitionsRevoked(partitions, rebalanceConsumer);
}
};
// Subscribe to get first assignment (no delays) and verify consumption
consumer.subscribe(List.of(topic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(topic));
var records = awaitNonEmptyRecords(consumer, tp, 0L);
assertEquals(numMessages, records.count());

// Subscribe to different topic. This will trigger the delayed revocation exceeding rebalance timeout and get fenced
consumer.subscribe(List.of(otherTopic), listener);
consumer.setRebalanceListener(listener);
consumer.subscribe(List.of(otherTopic));
ClientsTestUtils.pollUntilTrue(
consumer,
rebalanceTimeoutExceeded::get,
Expand Down
Loading