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 @@ -18,8 +18,9 @@

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.RebalanceConsumer;
import org.apache.kafka.clients.consumer.RebalanceListener;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.apache.kafka.common.serialization.Serde;
Expand Down Expand Up @@ -183,8 +184,8 @@ public void testRegexMatchesTopicsAWhenCreated() throws Exception {
public Consumer<byte[], byte[]> getConsumer(final Map<String, Object> config) {
return new KafkaConsumer<byte[], byte[]>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer()) {
@Override
public void subscribe(final Pattern topics, final ConsumerRebalanceListener listener) {
super.subscribe(topics, new TheConsumerRebalanceListener(assignedTopics, listener));
public void setRebalanceListener(final RebalanceListener listener) {
super.setRebalanceListener(new TheConsumerRebalanceListener(assignedTopics, listener));
}
};

Expand Down Expand Up @@ -282,8 +283,8 @@ public void shouldNotCrashIfPatternMatchesTopicHasNoData() throws Exception {
public Consumer<byte[], byte[]> getConsumer(final Map<String, Object> config) {
return new KafkaConsumer<>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer()) {
@Override
public void subscribe(final Pattern topics, final ConsumerRebalanceListener listener) {
super.subscribe(topics, new TheConsumerRebalanceListener(assignedTopics, listener));
public void setRebalanceListener(final RebalanceListener listener) {
super.setRebalanceListener(new TheConsumerRebalanceListener(assignedTopics, listener));
}
};
}
Expand Down Expand Up @@ -341,8 +342,8 @@ public void testRegexMatchesTopicsAWhenDeleted() throws Exception {
public Consumer<byte[], byte[]> getConsumer(final Map<String, Object> config) {
return new KafkaConsumer<byte[], byte[]>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer()) {
@Override
public void subscribe(final Pattern topics, final ConsumerRebalanceListener listener) {
super.subscribe(topics, new TheConsumerRebalanceListener(assignedTopics, listener));
public void setRebalanceListener(final RebalanceListener listener) {
super.setRebalanceListener(new TheConsumerRebalanceListener(assignedTopics, listener));
}
};
}
Expand Down Expand Up @@ -454,8 +455,8 @@ public void testMultipleConsumersCanReadFromPartitionedTopic() throws Exception
public Consumer<byte[], byte[]> getConsumer(final Map<String, Object> config) {
return new KafkaConsumer<byte[], byte[]>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer()) {
@Override
public void subscribe(final Pattern topics, final ConsumerRebalanceListener listener) {
super.subscribe(topics, new TheConsumerRebalanceListener(leaderAssignment, listener));
public void setRebalanceListener(final RebalanceListener listener) {
super.setRebalanceListener(new TheConsumerRebalanceListener(leaderAssignment, listener));
}
};

Expand All @@ -466,8 +467,8 @@ public void subscribe(final Pattern topics, final ConsumerRebalanceListener list
public Consumer<byte[], byte[]> getConsumer(final Map<String, Object> config) {
return new KafkaConsumer<byte[], byte[]>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer()) {
@Override
public void subscribe(final Pattern topics, final ConsumerRebalanceListener listener) {
super.subscribe(topics, new TheConsumerRebalanceListener(followerAssignment, listener));
public void setRebalanceListener(final RebalanceListener listener) {
super.setRebalanceListener(new TheConsumerRebalanceListener(followerAssignment, listener));
}
};

Expand Down Expand Up @@ -532,30 +533,30 @@ public void testNoMessagesSentExceptionFromOverlappingPatterns() throws Exceptio
assertThat(expectError.get(), is(true));
}

private static class TheConsumerRebalanceListener implements ConsumerRebalanceListener {
private static class TheConsumerRebalanceListener implements RebalanceListener {
private final List<String> assignedTopics;
private final ConsumerRebalanceListener listener;
private final RebalanceListener listener;

TheConsumerRebalanceListener(final List<String> assignedTopics, final ConsumerRebalanceListener listener) {
TheConsumerRebalanceListener(final List<String> assignedTopics, final RebalanceListener listener) {
this.assignedTopics = assignedTopics;
this.listener = listener;
}

@Override
public void onPartitionsRevoked(final Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(final Collection<TopicPartition> partitions, final RebalanceConsumer consumer) {
for (final TopicPartition partition : partitions) {
assignedTopics.remove(partition.topic());
}
listener.onPartitionsRevoked(partitions);
listener.onPartitionsRevoked(partitions, consumer);
}

@Override
public void onPartitionsAssigned(final Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(final Collection<TopicPartition> partitions, final RebalanceConsumer consumer) {
for (final TopicPartition partition : partitions) {
assignedTopics.add(partition.topic());
}
Collections.sort(assignedTopics);
listener.onPartitionsAssigned(partitions);
listener.onPartitionsAssigned(partitions, consumer);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,10 +24,10 @@
import org.apache.kafka.clients.admin.ListTopicsOptions;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.GroupProtocol;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.consumer.RebalanceListener;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
Expand Down Expand Up @@ -401,13 +401,12 @@ public KafkaConsumer<byte[], byte[]> createConsumerAndSubscribeTo(final Map<Stri
return createConsumerAndSubscribeTo(consumerProps, null, topics);
}

public KafkaConsumer<byte[], byte[]> createConsumerAndSubscribeTo(final Map<String, Object> consumerProps, final ConsumerRebalanceListener rebalanceListener, final String... topics) {
public KafkaConsumer<byte[], byte[]> createConsumerAndSubscribeTo(final Map<String, Object> consumerProps, final RebalanceListener rebalanceListener, final String... topics) {
final KafkaConsumer<byte[], byte[]> consumer = createConsumer(consumerProps);
if (rebalanceListener != null) {
consumer.subscribe(Arrays.asList(topics), rebalanceListener);
} else {
consumer.subscribe(Arrays.asList(topics));
consumer.setRebalanceListener(rebalanceListener);
}
consumer.subscribe(Arrays.asList(topics));
return consumer;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,12 +22,12 @@
import org.apache.kafka.clients.consumer.CloseOptions.GroupMembershipOperation;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.InvalidOffsetException;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.consumer.OffsetAndTimestamp;
import org.apache.kafka.clients.consumer.RebalanceListener;
import org.apache.kafka.clients.consumer.internals.AsyncKafkaConsumer;
import org.apache.kafka.clients.consumer.internals.AutoOffsetResetStrategy;
import org.apache.kafka.clients.consumer.internals.StreamsRebalanceData;
Expand Down Expand Up @@ -353,7 +353,7 @@ public boolean isStartingRunningOrPartitionAssigned() {
private final Optional<String> groupInstanceID;

private final ChangelogReader changelogReader;
private final ConsumerRebalanceListener rebalanceListener;
private final RebalanceListener rebalanceListener;
private final Optional<DefaultStreamsRebalanceListener> defaultStreamsRebalanceListener;
private final Consumer<byte[], byte[]> mainConsumer;
private final Consumer<byte[], byte[]> restoreConsumer;
Expand Down Expand Up @@ -1186,7 +1186,8 @@ private void subscribeConsumer() {
throw new IllegalArgumentException("Pattern subscription is not yet supported with the Streams rebalance " +
"protocol");
}
mainConsumer.subscribe(topologyMetadata.sourceTopicPattern(), rebalanceListener);
mainConsumer.setRebalanceListener(rebalanceListener);
mainConsumer.subscribe(topologyMetadata.sourceTopicPattern());
} else {
if (streamsRebalanceData.isPresent()) {
if (mainConsumer instanceof ConsumerWrapper) {
Expand All @@ -1201,7 +1202,8 @@ private void subscribeConsumer() {
);
}
} else {
mainConsumer.subscribe(topologyMetadata.allFullSourceTopicNames(), rebalanceListener);
mainConsumer.setRebalanceListener(rebalanceListener);
mainConsumer.subscribe(topologyMetadata.allFullSourceTopicNames());
}
}
}
Expand Down Expand Up @@ -2169,7 +2171,7 @@ int currentNumIterations() {
return numIterations;
}

ConsumerRebalanceListener rebalanceListener() {
RebalanceListener rebalanceListener() {
return rebalanceListener;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,8 @@
*/
package org.apache.kafka.streams.processor.internals;

import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.RebalanceConsumer;
import org.apache.kafka.clients.consumer.RebalanceListener;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.streams.errors.MissingSourceTopicException;
Expand All @@ -29,7 +30,7 @@
import java.util.Collection;
import java.util.concurrent.atomic.AtomicInteger;

public class StreamsRebalanceListener implements ConsumerRebalanceListener {
public class StreamsRebalanceListener implements RebalanceListener {

private final Time time;
private final TaskManager taskManager;
Expand All @@ -50,7 +51,7 @@ public class StreamsRebalanceListener implements ConsumerRebalanceListener {
}

@Override
public void onPartitionsAssigned(final Collection<TopicPartition> partitions) {
public void onPartitionsAssigned(final Collection<TopicPartition> partitions, final RebalanceConsumer consumer) {
// NB: all task management is already handled by:
// org.apache.kafka.streams.processor.internals.StreamsPartitionAssignor.onAssignment
if (assignmentErrorCode.get() == AssignorError.INCOMPLETE_SOURCE_TOPIC_METADATA.code()) {
Expand Down Expand Up @@ -81,7 +82,7 @@ public void onPartitionsAssigned(final Collection<TopicPartition> partitions) {
}

@Override
public void onPartitionsRevoked(final Collection<TopicPartition> partitions) {
public void onPartitionsRevoked(final Collection<TopicPartition> partitions, final RebalanceConsumer consumer) {
log.debug("Current state {}: revoked partitions {} because of consumer rebalance.\n" +
"\tcurrently assigned active tasks: {}\n" +
"\tcurrently assigned standby tasks: {}\n",
Expand All @@ -103,7 +104,7 @@ public void onPartitionsRevoked(final Collection<TopicPartition> partitions) {
}

@Override
public void onPartitionsLost(final Collection<TopicPartition> partitions) {
public void onPartitionsLost(final Collection<TopicPartition> partitions, final RebalanceConsumer consumer) {
log.info("at state {}: partitions {} lost due to missed rebalance.\n" +
"\tlost active tasks: {}\n" +
"\tlost assigned standby tasks: {}\n",
Expand Down
Loading
Loading