From aeaedeeee9f765b3d848e52cd2244036441cce71 Mon Sep 17 00:00:00 2001 From: Aditya Kousik Date: Thu, 6 Aug 2026 07:13:09 -0700 Subject: [PATCH] KAFKA-20684 [3/N]: Migrate tools to RebalanceListener Switches ConsumerPerfRebListener, VerifiableConsumer and TransactionalMessageCopier to RebalanceListener and registers them via setRebalanceListener. --- .../kafka/tools/ConsumerPerformance.java | 14 ++++++++------ .../tools/TransactionalMessageCopier.java | 10 ++++++---- .../apache/kafka/tools/VerifiableConsumer.java | 18 +++++++++++------- .../kafka/tools/ConsumerPerformanceTest.java | 16 ++++++++-------- 4 files changed, 33 insertions(+), 25 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/ConsumerPerformance.java b/tools/src/main/java/org/apache/kafka/tools/ConsumerPerformance.java index f38a20c41076e..451eda2f69952 100644 --- a/tools/src/main/java/org/apache/kafka/tools/ConsumerPerformance.java +++ b/tools/src/main/java/org/apache/kafka/tools/ConsumerPerformance.java @@ -18,10 +18,11 @@ 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.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.utils.Utils; @@ -140,10 +141,11 @@ private static void consume(Consumer consumer, SimpleDateFormat dateFormat = options.dateFormat(); ConsumerPerfRebListener listener = new ConsumerPerfRebListener(joinTimeMs, joinStartMs, joinTimeMsInSingleRound); + consumer.setRebalanceListener(listener); if (options.topic().isPresent()) { - consumer.subscribe(options.topic().get(), listener); + consumer.subscribe(options.topic().get()); } else { - consumer.subscribe(options.include().get(), listener); + consumer.subscribe(options.include().get()); } // now start the benchmark @@ -228,7 +230,7 @@ private static void printExtendedProgress(long bytesRead, fetchTimeMs, intervalMbPerSec, intervalRecordsPerSec); } - public static class ConsumerPerfRebListener implements ConsumerRebalanceListener { + public static class ConsumerPerfRebListener implements RebalanceListener { private final AtomicLong joinTimeMs; private final AtomicLong joinTimeMsInSingleRound; private final Collection assignedPartitions; @@ -242,7 +244,7 @@ public ConsumerPerfRebListener(AtomicLong joinTimeMs, long joinStartMs, AtomicLo } @Override - public void onPartitionsRevoked(Collection partitions) { + public void onPartitionsRevoked(Collection partitions, RebalanceConsumer consumer) { assignedPartitions.removeAll(partitions); if (assignedPartitions.isEmpty()) { joinStartMs = System.currentTimeMillis(); @@ -250,7 +252,7 @@ public void onPartitionsRevoked(Collection partitions) { } @Override - public void onPartitionsAssigned(Collection partitions) { + public void onPartitionsAssigned(Collection partitions, RebalanceConsumer consumer) { if (assignedPartitions.isEmpty()) { long elapsedMs = System.currentTimeMillis() - joinStartMs; joinTimeMs.addAndGet(elapsedMs); diff --git a/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java b/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java index 35a9a32fc47c6..01dcd6296abe1 100644 --- a/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java +++ b/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java @@ -18,11 +18,12 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerGroupMetadata; -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.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.clients.consumer.RebalanceConsumer; +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; @@ -319,13 +320,13 @@ public static void runEventLoop(Namespace parsedArgs) { final AtomicLong numMessagesProcessedSinceLastRebalance = new AtomicLong(0); final AtomicLong totalMessageProcessed = new AtomicLong(0); if (groupMode) { - consumer.subscribe(Set.of(topicName), new ConsumerRebalanceListener() { + consumer.setRebalanceListener(new RebalanceListener() { @Override - public void onPartitionsRevoked(Collection partitions) { + public void onPartitionsRevoked(Collection partitions, RebalanceConsumer rebalanceConsumer) { } @Override - public void onPartitionsAssigned(Collection partitions) { + public void onPartitionsAssigned(Collection partitions, RebalanceConsumer rebalanceConsumer) { remainingMessages.set(partitions.stream() .mapToLong(partition -> messagesRemaining(consumer, partition)).sum()); numMessagesProcessedSinceLastRebalance.set(0); @@ -339,6 +340,7 @@ public void onPartitionsAssigned(Collection partitions) { )); } }); + consumer.subscribe(Set.of(topicName)); } else { TopicPartition inputPartition = new TopicPartition(topicName, parsedArgs.getInt("inputPartition")); consumer.assign(Set.of(inputPartition)); diff --git a/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java b/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java index 242bd4992e1ac..41fd122aaec44 100644 --- a/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java +++ b/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java @@ -18,7 +18,6 @@ import org.apache.kafka.clients.consumer.CloseOptions; 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.GroupProtocol; @@ -26,6 +25,8 @@ import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.OffsetCommitCallback; import org.apache.kafka.clients.consumer.RangeAssignor; +import org.apache.kafka.clients.consumer.RebalanceConsumer; +import org.apache.kafka.clients.consumer.RebalanceListener; import org.apache.kafka.clients.consumer.RoundRobinAssignor; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.FencedInstanceIdException; @@ -75,9 +76,11 @@ * events are currently supported: * *
    - *
  • partitions_revoked: outputs the partitions revoked through {@link ConsumerRebalanceListener#onPartitionsRevoked(Collection)}. + *
  • partitions_revoked: outputs the partitions revoked through + * {@link RebalanceListener#onPartitionsRevoked(Collection, RebalanceConsumer)}. * See {@link org.apache.kafka.tools.VerifiableConsumer.PartitionsRevoked}.
  • - *
  • partitions_assigned: outputs the partitions assigned through {@link ConsumerRebalanceListener#onPartitionsAssigned(Collection)} + *
  • partitions_assigned: outputs the partitions assigned through + * {@link RebalanceListener#onPartitionsAssigned(Collection, RebalanceConsumer)} * See {@link org.apache.kafka.tools.VerifiableConsumer.PartitionsAssigned}.
  • *
  • records_consumed: contains a summary of records consumed in a single call to {@link KafkaConsumer#poll(Duration)}. * See {@link org.apache.kafka.tools.VerifiableConsumer.RecordsConsumed}.
  • @@ -91,7 +94,7 @@ * See {@link org.apache.kafka.tools.VerifiableConsumer.ShutdownComplete}. *
*/ -public class VerifiableConsumer implements Closeable, OffsetCommitCallback, ConsumerRebalanceListener { +public class VerifiableConsumer implements Closeable, OffsetCommitCallback, RebalanceListener { private static final Logger log = LoggerFactory.getLogger(VerifiableConsumer.class); @@ -201,12 +204,12 @@ public void onComplete(Map offsets, Exception } @Override - public void onPartitionsAssigned(Collection partitions) { + public void onPartitionsAssigned(Collection partitions, RebalanceConsumer rebalanceConsumer) { printJson(new PartitionsAssigned(partitions)); } @Override - public void onPartitionsRevoked(Collection partitions) { + public void onPartitionsRevoked(Collection partitions, RebalanceConsumer rebalanceConsumer) { printJson(new PartitionsRevoked(partitions)); } @@ -236,7 +239,8 @@ public void commitSync(Map offsets) { public void run() { try { printJson(new StartupComplete()); - consumer.subscribe(List.of(topic), this); + consumer.setRebalanceListener(this); + consumer.subscribe(List.of(topic)); while (!isFinished()) { ConsumerRecords records = consumer.poll(Duration.ofMillis(Long.MAX_VALUE)); diff --git a/tools/src/test/java/org/apache/kafka/tools/ConsumerPerformanceTest.java b/tools/src/test/java/org/apache/kafka/tools/ConsumerPerformanceTest.java index 7801bab0a78aa..f398320d17b34 100644 --- a/tools/src/test/java/org/apache/kafka/tools/ConsumerPerformanceTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/ConsumerPerformanceTest.java @@ -314,13 +314,13 @@ public void testConsumerListenerWithAllPartitionRevokedAndAssigned() throws Inte AtomicLong joinTimeMs = new AtomicLong(0); AtomicLong joinTimeMsInSingleRound = new AtomicLong(0); ConsumerPerformance.ConsumerPerfRebListener listener = new ConsumerPerformance.ConsumerPerfRebListener(joinTimeMs, 0, joinTimeMsInSingleRound); - listener.onPartitionsAssigned(Set.of(tp0)); + listener.onPartitionsAssigned(Set.of(tp0), null); long lastJoinTimeMs = joinTimeMs.get(); // All assigned partitions have been revoked. - listener.onPartitionsRevoked(Set.of(tp0)); + listener.onPartitionsRevoked(Set.of(tp0), null); Thread.sleep(100); - listener.onPartitionsAssigned(Set.of(tp1)); + listener.onPartitionsAssigned(Set.of(tp1), null); assertNotEquals(lastJoinTimeMs, joinTimeMs.get()); } @@ -333,13 +333,13 @@ public void testConsumerListenerWithPartialPartitionRevokedAndAssigned() throws AtomicLong joinTimeMs = new AtomicLong(0); AtomicLong joinTimeMsInSingleRound = new AtomicLong(0); ConsumerPerformance.ConsumerPerfRebListener listener = new ConsumerPerformance.ConsumerPerfRebListener(joinTimeMs, 0, joinTimeMsInSingleRound); - listener.onPartitionsAssigned(Set.of(tp0, tp1)); + listener.onPartitionsAssigned(Set.of(tp0, tp1), null); long lastJoinTimeMs = joinTimeMs.get(); // The assigned partitions were partially revoked. - listener.onPartitionsRevoked(Set.of(tp0)); + listener.onPartitionsRevoked(Set.of(tp0), null); Thread.sleep(100); - listener.onPartitionsAssigned(Set.of(tp0)); + listener.onPartitionsAssigned(Set.of(tp0), null); assertEquals(lastJoinTimeMs, joinTimeMs.get()); } @@ -352,11 +352,11 @@ public void testConsumerListenerWithoutPartitionRevoked() throws InterruptedExce AtomicLong joinTimeMs = new AtomicLong(0); AtomicLong joinTimeMsInSingleRound = new AtomicLong(0); ConsumerPerformance.ConsumerPerfRebListener listener = new ConsumerPerformance.ConsumerPerfRebListener(joinTimeMs, 0, joinTimeMsInSingleRound); - listener.onPartitionsAssigned(Set.of(tp0)); + listener.onPartitionsAssigned(Set.of(tp0), null); long lastJoinTimeMs = joinTimeMs.get(); Thread.sleep(100); - listener.onPartitionsAssigned(Set.of(tp1)); + listener.onPartitionsAssigned(Set.of(tp1), null); assertEquals(lastJoinTimeMs, joinTimeMs.get()); }