From 13c821ca2ef4f7b25050ae83403737629fadd1e2 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 14:04:23 +0800 Subject: [PATCH 01/16] init integrationTest --- .../group/ResetConsumerGroupOffsetTest.java | 35 +++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 5fb704cf53d35..e50d9e7d2f5ec 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.tools.consumer.group; import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewPartitionReassignment; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.GroupProtocol; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -26,7 +27,9 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.GroupState; +import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.TopicPartitionInfo; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -51,7 +54,9 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.Properties; +import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.function.Function; import java.util.function.Supplier; @@ -153,6 +158,36 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc } } + @ClusterTest(brokers = 5) + public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws Exception { + String topic = generateRandomTopic(); + String group = "new.group"; + String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", "topic:1"); + cluster.createTopic(topic, 3, (short) 2); + + try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args); + Admin admin = cluster.admin()) { + admin.alterPartitionReassignments(Map.of(new TopicPartition(topic, 2), + Optional.of(new NewPartitionReassignment(List.of(3, 4))))); + + System.err.println(admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions()); + + TestUtils.waitForCondition(() -> { + List partInfo = admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions(); + System.err.println(partInfo.get(2)); + return partInfo.get(2).replicas().stream().map(Node::id).collect(Collectors.toUnmodifiableSet()).equals(Set.of(3, 4)); + }, "aaaa"); + + cluster.shutdownBroker(3); + cluster.shutdownBroker(4); + TestUtils.waitForCondition(() -> !cluster.aliveBrokers().keySet().containsAll(Set.of(3, 4)), "aaaa"); + + Map resetOffsets = service.resetOffsets().get(group); + assertTrue(resetOffsets.isEmpty()); + assertTrue(committedOffsets(cluster, topic, group).isEmpty()); + } + } + @ClusterTest public void testResetOffsetsExistingTopic(ClusterInstance cluster) { String topic = generateRandomTopic(); From a0fdb86360c1781cccda655e16408e6dbd4c9a79 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 15:06:16 +0800 Subject: [PATCH 02/16] use delete dir to create offline partition --- .../java/org/apache/kafka/test/TestUtils.java | 5 +++ .../group/ResetConsumerGroupOffsetTest.java | 36 +++++++------------ 2 files changed, 17 insertions(+), 24 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/test/TestUtils.java b/clients/src/test/java/org/apache/kafka/test/TestUtils.java index 7748bbe15f099..3d63d84a26843 100644 --- a/clients/src/test/java/org/apache/kafka/test/TestUtils.java +++ b/clients/src/test/java/org/apache/kafka/test/TestUtils.java @@ -287,6 +287,11 @@ public static File tempRelativeDir(String root) { return tempDirectory(rootFile.toPath(), null); } + public static boolean deleteDir(String root) { + File rootFile = new File(root); + return rootFile.delete(); + } + /** * Create a temporary relative directory in the specified parent directory with the given prefix. * diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index e50d9e7d2f5ec..5864414520992 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -17,7 +17,6 @@ package org.apache.kafka.tools.consumer.group; import org.apache.kafka.clients.admin.Admin; -import org.apache.kafka.clients.admin.NewPartitionReassignment; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.GroupProtocol; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -27,9 +26,7 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.GroupState; -import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.TopicPartitionInfo; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -54,9 +51,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.Properties; -import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.function.Function; import java.util.function.Supplier; @@ -158,31 +153,24 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc } } - @ClusterTest(brokers = 5) - public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws Exception { + @ClusterTest(brokers = 3) + public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws InterruptedException { String topic = generateRandomTopic(); String group = "new.group"; - String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", "topic:1"); + String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":1"); cluster.createTopic(topic, 3, (short) 2); + TopicPartition offlinePartition = new TopicPartition(topic, 2); - try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args); - Admin admin = cluster.admin()) { - admin.alterPartitionReassignments(Map.of(new TopicPartition(topic, 2), - Optional.of(new NewPartitionReassignment(List.of(3, 4))))); - - System.err.println(admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions()); - - TestUtils.waitForCondition(() -> { - List partInfo = admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions(); - System.err.println(partInfo.get(2)); - return partInfo.get(2).replicas().stream().map(Node::id).collect(Collectors.toUnmodifiableSet()).equals(Set.of(3, 4)); - }, "aaaa"); - - cluster.shutdownBroker(3); - cluster.shutdownBroker(4); - TestUtils.waitForCondition(() -> !cluster.aliveBrokers().keySet().containsAll(Set.of(3, 4)), "aaaa"); + try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) { + cluster.brokers().forEach((id, broker) -> { + if (broker.replicaManager().onlinePartition(offlinePartition).isDefined()) { + System.err.println("KKKK" + broker.logManager().getLog(offlinePartition, false).get().dir()); + TestUtils.deleteDir(broker.logManager().getLog(offlinePartition, false).get().dir().getAbsolutePath()); + } + }); Map resetOffsets = service.resetOffsets().get(group); + System.err.println("ZZZZ " + resetOffsets); assertTrue(resetOffsets.isEmpty()); assertTrue(committedOffsets(cluster, topic, group).isEmpty()); } From f7a4b03255e85d404d664c4283a4a478ded61672 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 15:20:01 +0800 Subject: [PATCH 03/16] Revert "use delete dir to create offline partition" This reverts commit a0fdb86360c1781cccda655e16408e6dbd4c9a79. --- .../java/org/apache/kafka/test/TestUtils.java | 5 --- .../group/ResetConsumerGroupOffsetTest.java | 36 ++++++++++++------- 2 files changed, 24 insertions(+), 17 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/test/TestUtils.java b/clients/src/test/java/org/apache/kafka/test/TestUtils.java index 3d63d84a26843..7748bbe15f099 100644 --- a/clients/src/test/java/org/apache/kafka/test/TestUtils.java +++ b/clients/src/test/java/org/apache/kafka/test/TestUtils.java @@ -287,11 +287,6 @@ public static File tempRelativeDir(String root) { return tempDirectory(rootFile.toPath(), null); } - public static boolean deleteDir(String root) { - File rootFile = new File(root); - return rootFile.delete(); - } - /** * Create a temporary relative directory in the specified parent directory with the given prefix. * diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 5864414520992..e50d9e7d2f5ec 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.tools.consumer.group; import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewPartitionReassignment; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.GroupProtocol; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -26,7 +27,9 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.GroupState; +import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.TopicPartitionInfo; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -51,7 +54,9 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.Properties; +import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.function.Function; import java.util.function.Supplier; @@ -153,24 +158,31 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc } } - @ClusterTest(brokers = 3) - public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws InterruptedException { + @ClusterTest(brokers = 5) + public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws Exception { String topic = generateRandomTopic(); String group = "new.group"; - String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":1"); + String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", "topic:1"); cluster.createTopic(topic, 3, (short) 2); - TopicPartition offlinePartition = new TopicPartition(topic, 2); - try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) { - cluster.brokers().forEach((id, broker) -> { - if (broker.replicaManager().onlinePartition(offlinePartition).isDefined()) { - System.err.println("KKKK" + broker.logManager().getLog(offlinePartition, false).get().dir()); - TestUtils.deleteDir(broker.logManager().getLog(offlinePartition, false).get().dir().getAbsolutePath()); - } - }); + try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args); + Admin admin = cluster.admin()) { + admin.alterPartitionReassignments(Map.of(new TopicPartition(topic, 2), + Optional.of(new NewPartitionReassignment(List.of(3, 4))))); + + System.err.println(admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions()); + + TestUtils.waitForCondition(() -> { + List partInfo = admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions(); + System.err.println(partInfo.get(2)); + return partInfo.get(2).replicas().stream().map(Node::id).collect(Collectors.toUnmodifiableSet()).equals(Set.of(3, 4)); + }, "aaaa"); + + cluster.shutdownBroker(3); + cluster.shutdownBroker(4); + TestUtils.waitForCondition(() -> !cluster.aliveBrokers().keySet().containsAll(Set.of(3, 4)), "aaaa"); Map resetOffsets = service.resetOffsets().get(group); - System.err.println("ZZZZ " + resetOffsets); assertTrue(resetOffsets.isEmpty()); assertTrue(committedOffsets(cluster, topic, group).isEmpty()); } From a18d896011efc8848e4b093840f958d1adabc910 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 15:28:11 +0800 Subject: [PATCH 04/16] add integration test --- .../tools/consumer/group/ResetConsumerGroupOffsetTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index e50d9e7d2f5ec..931055d6eb115 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -162,7 +162,7 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws Exception { String topic = generateRandomTopic(); String group = "new.group"; - String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", "topic:1"); + String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":1"); cluster.createTopic(topic, 3, (short) 2); try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args); From c2a24c9ac554251aa4c50780ba9548ec3ee0dd93 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 15:46:49 +0800 Subject: [PATCH 05/16] fix --- .../tools/consumer/group/ConsumerGroupCommand.java | 4 ++-- .../consumer/group/ResetConsumerGroupOffsetTest.java | 11 ++++------- 2 files changed, 6 insertions(+), 9 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java index 2cfffb9fe7881..bbaee6d782e78 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java @@ -653,7 +653,7 @@ private List filterNoneLeaderPartitions(Collection entry.getValue().partitions().stream() - .filter(partitionInfo -> partitionInfo.leader() == null) + .filter(partitionInfo -> topicPartitions.contains(new TopicPartition(entry.getKey(), partitionInfo.partition())) && partitionInfo.leader() == null) .map(partitionInfo -> new TopicPartition(entry.getKey(), partitionInfo.partition()))) .toList(); } catch (Exception e) { @@ -1054,7 +1054,7 @@ private List filterNonExistentPartitions(Collection new TopicPartition(entry.getKey(), partitionInfo.partition()))) .toList(); - return topicPartitions.stream().filter(element -> !existPartitions.contains(element)).toList(); + return topicPartitions.stream().filter(tp -> !existPartitions.contains(tp)).toList(); } catch (InterruptedException | ExecutionException e) { throw new RuntimeException(e); } diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 931055d6eb115..d50995f71d7cb 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -170,21 +170,18 @@ public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws admin.alterPartitionReassignments(Map.of(new TopicPartition(topic, 2), Optional.of(new NewPartitionReassignment(List.of(3, 4))))); - System.err.println(admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions()); - TestUtils.waitForCondition(() -> { List partInfo = admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions(); - System.err.println(partInfo.get(2)); return partInfo.get(2).replicas().stream().map(Node::id).collect(Collectors.toUnmodifiableSet()).equals(Set.of(3, 4)); - }, "aaaa"); + }, "Partitions did not complete reassignment to the target broker(s) within the expected time."); cluster.shutdownBroker(3); cluster.shutdownBroker(4); - TestUtils.waitForCondition(() -> !cluster.aliveBrokers().keySet().containsAll(Set.of(3, 4)), "aaaa"); + TestUtils.waitForCondition(() -> !cluster.aliveBrokers().keySet().containsAll(Set.of(3, 4)), + "The cluster did not shut down the broker within the expected time."); Map resetOffsets = service.resetOffsets().get(group); - assertTrue(resetOffsets.isEmpty()); - assertTrue(committedOffsets(cluster, topic, group).isEmpty()); + assertEquals(Set.of(new TopicPartition(topic, 1)), resetOffsets.keySet()); } } From b0e292fb2b38325777fbda322719fb09490961d5 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 16:09:53 +0800 Subject: [PATCH 06/16] address indentation comment --- .../tools/consumer/group/ResetConsumerGroupOffsetTest.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index d50995f71d7cb..e29b305d064a3 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -165,8 +165,7 @@ public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":1"); cluster.createTopic(topic, 3, (short) 2); - try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args); - Admin admin = cluster.admin()) { + try (Admin admin = cluster.admin(); ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) { admin.alterPartitionReassignments(Map.of(new TopicPartition(topic, 2), Optional.of(new NewPartitionReassignment(List.of(3, 4))))); From 391326a1438b4937be8dafbd45bf229ef0535c7b Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 17:47:01 +0800 Subject: [PATCH 07/16] fix flaky --- .../group/ResetConsumerGroupOffsetTest.java | 14 ++------------ 1 file changed, 2 insertions(+), 12 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index e29b305d064a3..8c1f90dc023c4 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -17,7 +17,6 @@ package org.apache.kafka.tools.consumer.group; import org.apache.kafka.clients.admin.Admin; -import org.apache.kafka.clients.admin.NewPartitionReassignment; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.GroupProtocol; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -27,9 +26,7 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.GroupState; -import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.TopicPartitionInfo; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -54,7 +51,6 @@ import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.Properties; import java.util.Set; import java.util.concurrent.ExecutionException; @@ -163,16 +159,10 @@ public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws String topic = generateRandomTopic(); String group = "new.group"; String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":1"); - cluster.createTopic(topic, 3, (short) 2); try (Admin admin = cluster.admin(); ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) { - admin.alterPartitionReassignments(Map.of(new TopicPartition(topic, 2), - Optional.of(new NewPartitionReassignment(List.of(3, 4))))); - - TestUtils.waitForCondition(() -> { - List partInfo = admin.describeTopics(List.of(topic)).topicNameValues().get(topic).get().partitions(); - return partInfo.get(2).replicas().stream().map(Node::id).collect(Collectors.toUnmodifiableSet()).equals(Set.of(3, 4)); - }, "Partitions did not complete reassignment to the target broker(s) within the expected time."); + admin.createTopics(List.of(new NewTopic(topic, Map.of(0, List.of(0, 1), 1, List.of(1, 2), 2, List.of(3, 4))))); + cluster.waitTopicCreation(topic, 3); cluster.shutdownBroker(3); cluster.shutdownBroker(4); From e1cfa54b65a1d5d42a95573e153420d03d540b42 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 21:33:09 +0800 Subject: [PATCH 08/16] address Yunyung comment --- .../group/ResetConsumerGroupOffsetTest.java | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 8c1f90dc023c4..4d4a282313f3b 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -154,23 +154,27 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc } } - @ClusterTest(brokers = 5) - public void testResetOffsetsWithOfflinePartition(ClusterInstance cluster) throws Exception { + @ClusterTest( + brokers = 3, + serverProperties = { + @ClusterConfigProperty(key = OFFSETS_TOPIC_REPLICATION_FACTOR_CONFIG, value = "2") + } + ) + public void testResetOffsetsWithOfflinePartitionNotInResetTarget(ClusterInstance cluster) throws Exception { String topic = generateRandomTopic(); String group = "new.group"; - String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":1"); + String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":0"); try (Admin admin = cluster.admin(); ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) { - admin.createTopics(List.of(new NewTopic(topic, Map.of(0, List.of(0, 1), 1, List.of(1, 2), 2, List.of(3, 4))))); - cluster.waitTopicCreation(topic, 3); + admin.createTopics(List.of(new NewTopic(topic, Map.of(0, List.of(0), 1, List.of(1))))); + cluster.waitTopicCreation(topic, 2); - cluster.shutdownBroker(3); - cluster.shutdownBroker(4); - TestUtils.waitForCondition(() -> !cluster.aliveBrokers().keySet().containsAll(Set.of(3, 4)), + cluster.shutdownBroker(1); + TestUtils.waitForCondition(() -> !cluster.aliveBrokers().containsKey(1), "The cluster did not shut down the broker within the expected time."); Map resetOffsets = service.resetOffsets().get(group); - assertEquals(Set.of(new TopicPartition(topic, 1)), resetOffsets.keySet()); + assertEquals(Set.of(new TopicPartition(topic, 0)), resetOffsets.keySet()); } } From e13bc80464c0c44c3327083298667856d8d56e8e Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 22:16:33 +0800 Subject: [PATCH 09/16] improve performance --- .../apache/kafka/tools/consumer/group/ConsumerGroupCommand.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java index bbaee6d782e78..62e2e486bb5d5 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java @@ -653,7 +653,7 @@ private List filterNoneLeaderPartitions(Collection entry.getValue().partitions().stream() - .filter(partitionInfo -> topicPartitions.contains(new TopicPartition(entry.getKey(), partitionInfo.partition())) && partitionInfo.leader() == null) + .filter(partitionInfo -> partitionInfo.leader() == null && topicPartitions.contains(new TopicPartition(entry.getKey(), partitionInfo.partition()))) .map(partitionInfo -> new TopicPartition(entry.getKey(), partitionInfo.partition()))) .toList(); } catch (Exception e) { From 05b504b07aef8478cc2e3ed44907a5728be68ac1 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 22:20:15 +0800 Subject: [PATCH 10/16] wait shutdown for testResetOffsetsWithPartitionNoneLeader --- .../tools/consumer/group/ResetConsumerGroupOffsetTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 4d4a282313f3b..c511830c3eab1 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -702,6 +702,8 @@ public void testResetOffsetsWithPartitionNoneLeader(ClusterInstance cluster) thr assertDoesNotThrow(() -> resetOffsets(service)); // shutdown a broker to make some partitions missing leader cluster.shutdownBroker(0); + TestUtils.waitForCondition(() -> !cluster.aliveBrokers().containsKey(0), + "The cluster did not shut down the broker within the expected time."); assertThrows(LeaderNotAvailableException.class, () -> resetOffsets(service)); } } From 7cd07038a8a9e61ef649a281078d887aeb6f5619 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 22:27:08 +0800 Subject: [PATCH 11/16] Revert "wait shutdown for testResetOffsetsWithPartitionNoneLeader" This reverts commit 05b504b07aef8478cc2e3ed44907a5728be68ac1. --- .../tools/consumer/group/ResetConsumerGroupOffsetTest.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index c511830c3eab1..4d4a282313f3b 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -702,8 +702,6 @@ public void testResetOffsetsWithPartitionNoneLeader(ClusterInstance cluster) thr assertDoesNotThrow(() -> resetOffsets(service)); // shutdown a broker to make some partitions missing leader cluster.shutdownBroker(0); - TestUtils.waitForCondition(() -> !cluster.aliveBrokers().containsKey(0), - "The cluster did not shut down the broker within the expected time."); assertThrows(LeaderNotAvailableException.class, () -> resetOffsets(service)); } } From cbe8138151a702e363fe422e73f08904757491ed Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Thu, 24 Jul 2025 23:30:39 +0800 Subject: [PATCH 12/16] remove wait --- .../tools/consumer/group/ResetConsumerGroupOffsetTest.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 4d4a282313f3b..7460c6a3282fe 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -170,8 +170,6 @@ public void testResetOffsetsWithOfflinePartitionNotInResetTarget(ClusterInstance cluster.waitTopicCreation(topic, 2); cluster.shutdownBroker(1); - TestUtils.waitForCondition(() -> !cluster.aliveBrokers().containsKey(1), - "The cluster did not shut down the broker within the expected time."); Map resetOffsets = service.resetOffsets().get(group); assertEquals(Set.of(new TopicPartition(topic, 0)), resetOffsets.keySet()); From 3406b2299296f938f8142d4235b711a97e977277 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Fri, 25 Jul 2025 07:44:38 +0800 Subject: [PATCH 13/16] address chia comments --- .../kafka/tools/consumer/group/ConsumerGroupCommand.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java index 62e2e486bb5d5..e0a9728f9a389 100644 --- a/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/consumer/group/ConsumerGroupCommand.java @@ -653,8 +653,9 @@ private List filterNoneLeaderPartitions(Collection entry.getValue().partitions().stream() - .filter(partitionInfo -> partitionInfo.leader() == null && topicPartitions.contains(new TopicPartition(entry.getKey(), partitionInfo.partition()))) + .filter(partitionInfo -> partitionInfo.leader() == null) .map(partitionInfo -> new TopicPartition(entry.getKey(), partitionInfo.partition()))) + .filter(topicPartitions::contains) .toList(); } catch (Exception e) { throw new RuntimeException(e); From b843e4e37f34f1234dd985f4b7150b1446764b7f Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Fri, 25 Jul 2025 16:41:36 +0800 Subject: [PATCH 14/16] set min_in_sync --- .../tools/consumer/group/ResetConsumerGroupOffsetTest.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 7460c6a3282fe..a55668ef70de6 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -27,6 +27,7 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.GroupState; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.TopicConfig; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -155,9 +156,10 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc } @ClusterTest( - brokers = 3, + brokers = 2, serverProperties = { - @ClusterConfigProperty(key = OFFSETS_TOPIC_REPLICATION_FACTOR_CONFIG, value = "2") + @ClusterConfigProperty(key = OFFSETS_TOPIC_REPLICATION_FACTOR_CONFIG, value = "2"), + @ClusterConfigProperty(key = TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, value = "1") } ) public void testResetOffsetsWithOfflinePartitionNotInResetTarget(ClusterInstance cluster) throws Exception { From 1a9e642b714b4a17f20716d3e0a282c9d4991279 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Wed, 30 Jul 2025 18:07:21 +0800 Subject: [PATCH 15/16] address chia comments --- .../kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index a55668ef70de6..2b4e2f17b7f1c 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -159,7 +159,6 @@ public void testResetOffsetsNotExistingGroup(ClusterInstance cluster) throws Exc brokers = 2, serverProperties = { @ClusterConfigProperty(key = OFFSETS_TOPIC_REPLICATION_FACTOR_CONFIG, value = "2"), - @ClusterConfigProperty(key = TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, value = "1") } ) public void testResetOffsetsWithOfflinePartitionNotInResetTarget(ClusterInstance cluster) throws Exception { From cda748f51bbaa39f7a8b3ea9f25b9add13d187c8 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Wed, 30 Jul 2025 20:09:22 +0800 Subject: [PATCH 16/16] fix build --- .../kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java index 2b4e2f17b7f1c..5bf9da0c3708e 100644 --- a/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java @@ -27,7 +27,6 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.GroupState; import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.config.TopicConfig; import org.apache.kafka.common.errors.LeaderNotAvailableException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; import org.apache.kafka.common.serialization.ByteArraySerializer;