From 4c4b8b08d4efa46cc1586b802f40aa42d8b062e4 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sat, 3 Aug 2024 01:27:59 +0800 Subject: [PATCH 01/10] KAFKA-16154: Broker returns offset for LATEST_TIERED_TIMESTAMP --- .../kafka/clients/admin/KafkaAdminClient.java | 4 +++ .../kafka/clients/admin/OffsetSpec.java | 19 ++++++++++++ .../apache/kafka/tools/GetOffsetShell.java | 8 +++++ .../kafka/tools/GetOffsetShellTest.java | 31 +++++++++++++++++++ 4 files changed, 62 insertions(+) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 640a08a3786a7..8eb7fb4e8c05e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -4860,6 +4860,10 @@ private static long getOffsetFromSpec(OffsetSpec offsetSpec) { return ListOffsetsRequest.EARLIEST_TIMESTAMP; } else if (offsetSpec instanceof OffsetSpec.MaxTimestampSpec) { return ListOffsetsRequest.MAX_TIMESTAMP; + } else if (offsetSpec instanceof OffsetSpec.EarliestLocalSpec) { + return ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP; + } else if (offsetSpec instanceof OffsetSpec.LatestTierSpec) { + return ListOffsetsRequest.LATEST_TIERED_TIMESTAMP; } return ListOffsetsRequest.LATEST_TIMESTAMP; } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java b/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java index dcf90452c55e7..5b2fbb3e2e950 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java @@ -26,6 +26,8 @@ public class OffsetSpec { public static class EarliestSpec extends OffsetSpec { } public static class LatestSpec extends OffsetSpec { } public static class MaxTimestampSpec extends OffsetSpec { } + public static class EarliestLocalSpec extends OffsetSpec { } + public static class LatestTierSpec extends OffsetSpec { } public static class TimestampSpec extends OffsetSpec { private final long timestamp; @@ -70,4 +72,21 @@ public static OffsetSpec maxTimestamp() { return new MaxTimestampSpec(); } + /** + * Used to retrieve the offset with the local log start offset, + * log start offset is the offset of a log above which reads are guaranteed to be served + * from the disk of the leader broker, when Tiered Storage is not enabled, it behaves the same + * as the earliest timestamp + */ + public static OffsetSpec earliestLocalSpec() { + return new EarliestLocalSpec(); + } + + /** + * Used to retrieve the offset with the highest offset of data stored in remote storage, + * and when Tiered Storage is not enabled, we won't return any offset (i.e. Unknown offset) + */ + public static OffsetSpec latestTierSpec() { + return new LatestTierSpec(); + } } diff --git a/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java b/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java index 4ba0f6c3e3c4d..9b1ebe9cfe90d 100644 --- a/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java +++ b/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java @@ -283,6 +283,10 @@ private OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseExce return OffsetSpec.latest(); case "max-timestamp": return OffsetSpec.maxTimestamp(); + case "earliest-local": + return OffsetSpec.earliestLocalSpec(); + case "latest-tiered": + return OffsetSpec.latestTierSpec(); default: long timestamp; @@ -299,6 +303,10 @@ private OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseExce return OffsetSpec.latest(); } else if (timestamp == ListOffsetsRequest.MAX_TIMESTAMP) { return OffsetSpec.maxTimestamp(); + } else if (timestamp == ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP) { + return OffsetSpec.earliestLocalSpec(); + } else if (timestamp == ListOffsetsRequest.LATEST_TIERED_TIMESTAMP) { + return OffsetSpec.latestTierSpec(); } else { return OffsetSpec.forTimestamp(timestamp); } diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index 2d588c6025734..dfe17e9252d00 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -41,6 +41,7 @@ import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.Objects; import java.util.Properties; @@ -274,6 +275,36 @@ public void testGetOffsetsByMaxTimestamp() { } } + @ClusterTest + public void testGetOffsetsByEarliestLocalSpec() { + setUp(); + + for (String time : new String[] {"-4", "earliest-local"}) { + List offsets = executeAndParse("--topic-partitions", "topic.*:0", "--time", time); + // as remote log not enabled, the result should be the same as earliest offsetspec + List expected = Arrays.asList( + new Row("topic1", 0, 0L), + new Row("topic2", 0, 0L), + new Row("topic3", 0, 0L), + new Row("topic4", 0, 0L) + ); + + assertEquals(expected, offsets); + } + } + + @ClusterTest + public void testGetOffsetsByLatestTieredSpec() { + setUp(); + + for (String time : new String[] {"-5", "latest-tiered"}) { + List offsets = executeAndParse("--topic-partitions", "topic.*:0", "--time", time); + // as remote log not enabled, broker return unknown offset for each topic partition and these + // unknown offsets are ignored by GetOffsetShell hence we have empty result here. + assertEquals(Collections.emptyList(), offsets); + } + } + @ClusterTest public void testGetOffsetsByTimestamp() { setUp(); From ceabb10aed79edf2732eb81294ad8cda54eeee66 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sat, 3 Aug 2024 09:38:36 +0800 Subject: [PATCH 02/10] Address comments --- .../org/apache/kafka/clients/admin/OffsetSpec.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java b/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java index 5b2fbb3e2e950..f2f8f6604022a 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java @@ -73,18 +73,20 @@ public static OffsetSpec maxTimestamp() { } /** - * Used to retrieve the offset with the local log start offset, - * log start offset is the offset of a log above which reads are guaranteed to be served - * from the disk of the leader broker, when Tiered Storage is not enabled, it behaves the same - * as the earliest timestamp + * Used to retrieve the local log start offset. + * Local log start offset is the offset of a log above which reads + * are guaranteed to be served from the disk of the leader broker. + *
+ * Note: When tiered Storage is not enabled, it behaves the same as retrieving the earliest timestamp offset. */ public static OffsetSpec earliestLocalSpec() { return new EarliestLocalSpec(); } /** - * Used to retrieve the offset with the highest offset of data stored in remote storage, - * and when Tiered Storage is not enabled, we won't return any offset (i.e. Unknown offset) + * Used to retrieve the highest offset of data stored in remote storage. + *
+ * Note: When tiered storage is not enabled, we will return unknown offset. */ public static OffsetSpec latestTierSpec() { return new LatestTierSpec(); From 55fb64886c131f0100cc8e2e269157cb6506e5f6 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sat, 3 Aug 2024 16:34:36 +0800 Subject: [PATCH 03/10] Add test --- build.gradle | 1 + checkstyle/import-control.xml | 2 + .../kafka/tools/GetOffsetShellTest.java | 117 ++++++++++++++---- 3 files changed, 95 insertions(+), 25 deletions(-) diff --git a/build.gradle b/build.gradle index 51f9659e587b0..becac56bf882b 100644 --- a/build.gradle +++ b/build.gradle @@ -2141,6 +2141,7 @@ project(':tools') { testImplementation project(':connect:runtime') testImplementation project(':connect:runtime').sourceSets.test.output testImplementation project(':storage:storage-api').sourceSets.main.output + testImplementation project(':storage').sourceSets.test.output testImplementation libs.junitJupiter testImplementation libs.mockitoCore testImplementation libs.mockitoJunitJupiter // supports MockitoExtension diff --git a/checkstyle/import-control.xml b/checkstyle/import-control.xml index a5784ef935cd1..3f8212f997637 100644 --- a/checkstyle/import-control.xml +++ b/checkstyle/import-control.xml @@ -284,6 +284,8 @@ + + diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index dfe17e9252d00..84a680efb472d 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -17,8 +17,10 @@ package org.apache.kafka.tools; +import kafka.test.ClusterConfig; import kafka.test.ClusterInstance; import kafka.test.annotation.ClusterConfigProperty; +import kafka.test.annotation.ClusterTemplate; import kafka.test.annotation.ClusterTest; import kafka.test.annotation.ClusterTestDefaults; import kafka.test.junit.ClusterTestExtensions; @@ -31,10 +33,16 @@ import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.config.TopicConfig; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.common.utils.AppInfoParser; import org.apache.kafka.common.utils.Exit; +import org.apache.kafka.server.config.ServerLogConfigs; +import org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig; +import org.apache.kafka.server.log.remote.storage.LocalTieredStorage; +import org.apache.kafka.server.log.remote.storage.RemoteLogManagerConfig; +import org.apache.kafka.test.TestUtils; import org.junit.jupiter.api.extension.ExtendWith; @@ -42,13 +50,19 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Properties; +import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.IntStream; import java.util.stream.Stream; +import static kafka.test.annotation.Type.CO_KRAFT; +import static kafka.test.annotation.Type.KRAFT; +import static kafka.test.annotation.Type.ZK; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -71,15 +85,36 @@ private String getTopicName(int i) { return "topic" + i; } + private String getRemoteLogStorageEnabledTopicName(int i) { + return "topicRLS" + i; + } + private void setUp() { + setupTopics(this::getTopicName, Collections.emptyMap()); + sendProducerRecords(this::getTopicName); + } + + private void setUpRemoteLogTopics() { + Map rlsConfigs = new HashMap<>(); + rlsConfigs.put(TopicConfig.REMOTE_LOG_STORAGE_ENABLE_CONFIG, "true"); + rlsConfigs.put(TopicConfig.LOCAL_LOG_RETENTION_BYTES_CONFIG, "1"); + rlsConfigs.put(TopicConfig.SEGMENT_BYTES_CONFIG, "100"); + setupTopics(this::getRemoteLogStorageEnabledTopicName, rlsConfigs); + sendProducerRecords(this::getRemoteLogStorageEnabledTopicName); + } + + private void setupTopics(Function topicName, Map configs) { try (Admin admin = cluster.createAdminClient()) { List topics = new ArrayList<>(); - IntStream.range(0, topicCount + 1).forEach(i -> topics.add(new NewTopic(getTopicName(i), i, (short) 1))); + IntStream.range(0, topicCount + 1).forEach(i -> + topics.add(new NewTopic(topicName.apply(i), i, (short) 1).configs(configs))); admin.createTopics(topics); } + } + private void sendProducerRecords(Function topicName) { Properties props = new Properties(); props.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, cluster.bootstrapServers()); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); @@ -87,15 +122,38 @@ private void setUp() { try (KafkaProducer producer = new KafkaProducer<>(props)) { IntStream.range(0, topicCount + 1) - .forEach(i -> IntStream.range(0, i * i) - .forEach(msgCount -> { - assertDoesNotThrow(() -> producer.send( - new ProducerRecord<>(getTopicName(i), msgCount % i, null, "val" + msgCount)).get()); - }) - ); + .forEach(i -> IntStream.range(0, i * i) + .forEach(msgCount -> assertDoesNotThrow(() -> producer.send( + new ProducerRecord<>(topicName.apply(i), msgCount % i, null, "val" + msgCount)).get()))); } } + private static List withRemoteStorage() { + Map serverProperties = new HashMap<>(); + serverProperties.put(RemoteLogManagerConfig.DEFAULT_REMOTE_LOG_METADATA_MANAGER_CONFIG_PREFIX + TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP, "1"); + serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_STORAGE_SYSTEM_ENABLE_PROP, "true"); + serverProperties.put(RemoteLogManagerConfig.REMOTE_STORAGE_MANAGER_CLASS_NAME_PROP, LocalTieredStorage.class.getName()); + serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_MANAGER_TASK_INTERVAL_MS_PROP, "1000"); + serverProperties.put(ServerLogConfigs.LOG_CLEANUP_INTERVAL_MS_CONFIG, "1000"); + serverProperties.put(ServerLogConfigs.LOG_INITIAL_TASK_DELAY_MS_CONFIG, "100"); + + Map zkProperties = new HashMap<>(serverProperties); + zkProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "PLAINTEXT"); + + Map raftProperties = new HashMap<>(serverProperties); + raftProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "EXTERNAL"); + + return Arrays.asList( + ClusterConfig.defaultBuilder() + .setTypes(Collections.singleton(ZK)) + .setServerProperties(zkProperties) + .build(), + ClusterConfig.defaultBuilder() + .setTypes(Stream.of(KRAFT, CO_KRAFT).collect(Collectors.toSet())) + .setServerProperties(raftProperties) + .build()); + } + private void createConsumerAndPoll() { Properties props = new Properties(); props.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, cluster.bootstrapServers()); @@ -275,33 +333,42 @@ public void testGetOffsetsByMaxTimestamp() { } } - @ClusterTest - public void testGetOffsetsByEarliestLocalSpec() { - setUp(); + @ClusterTemplate("withRemoteStorage") + public void testGetOffsetsByEarliestLocalSpec() throws InterruptedException { + setUpRemoteLogTopics(); for (String time : new String[] {"-4", "earliest-local"}) { - List offsets = executeAndParse("--topic-partitions", "topic.*:0", "--time", time); - // as remote log not enabled, the result should be the same as earliest offsetspec - List expected = Arrays.asList( - new Row("topic1", 0, 0L), - new Row("topic2", 0, 0L), - new Row("topic3", 0, 0L), - new Row("topic4", 0, 0L) - ); - - assertEquals(expected, offsets); + TestUtils.waitForCondition(() -> + Arrays.asList( + new Row("topicRLS1", 0, 0L), + new Row("topicRLS2", 0, 1L), + new Row("topicRLS3", 0, 2L), + new Row("topicRLS4", 0, 3L)) + .equals(executeAndParse("--topic-partitions", "topicRLS.*:0", "--time", time)), + "testGetOffsetsByEarliestLocalSpec result not match"); } } - @ClusterTest - public void testGetOffsetsByLatestTieredSpec() { + @ClusterTemplate("withRemoteStorage") + public void testGetOffsetsByLatestTieredSpec() throws InterruptedException { setUp(); + setUpRemoteLogTopics(); for (String time : new String[] {"-5", "latest-tiered"}) { - List offsets = executeAndParse("--topic-partitions", "topic.*:0", "--time", time); + // test topics disable remote log storage // as remote log not enabled, broker return unknown offset for each topic partition and these // unknown offsets are ignored by GetOffsetShell hence we have empty result here. - assertEquals(Collections.emptyList(), offsets); + assertEquals(Collections.emptyList(), + executeAndParse("--topic-partitions", "topic\\d+:0", "--time", time)); + + // test topics enable remote log storage + TestUtils.waitForCondition(() -> + Arrays.asList( + new Row("topicRLS2", 0, 0L), + new Row("topicRLS3", 0, 1L), + new Row("topicRLS4", 0, 2L)) + .equals(executeAndParse("--topic-partitions", "topicRLS.*:0", "--time", time)), + "testGetOffsetsByLatestTieredSpec result not match"); } } @@ -447,7 +514,7 @@ private List expectedOffsetsForTopic(int i) { private List executeAndParse(String... args) { String out = ToolsTestUtils.captureStandardOut(() -> GetOffsetShell.mainNoExit(addBootstrapServer(args))); - + System.out.println(out); return Arrays.stream(out.split(System.lineSeparator())) .map(i -> i.split(":")) .filter(i -> i.length >= 2) From 934b8f31e9ffe3bb83aa263b023d6693074d11b5 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sat, 3 Aug 2024 16:47:46 +0800 Subject: [PATCH 04/10] Remove debug log --- .../test/java/org/apache/kafka/tools/GetOffsetShellTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index 84a680efb472d..1ea2d1a493f01 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -514,7 +514,7 @@ private List expectedOffsetsForTopic(int i) { private List executeAndParse(String... args) { String out = ToolsTestUtils.captureStandardOut(() -> GetOffsetShell.mainNoExit(addBootstrapServer(args))); - System.out.println(out); + return Arrays.stream(out.split(System.lineSeparator())) .map(i -> i.split(":")) .filter(i -> i.length >= 2) From 87b5213dfbeed347e5246d96e2dbaa14d63b0cf4 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sun, 4 Aug 2024 01:04:07 +0800 Subject: [PATCH 05/10] Address comments --- .../kafka/tools/GetOffsetShellTest.java | 21 ++++++++++++------- 1 file changed, 13 insertions(+), 8 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index 1ea2d1a493f01..831baa8ab7224 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -95,6 +95,11 @@ private void setUp() { } private void setUpRemoteLogTopics() { + // In this method, we'll create 4 topics and produce records to the log like this: + // topicRLS1 -> 1 segment + // topicRLS2 -> 2 segments (1 local log segment + 1 segment in the remote storage) + // topicRLS3 -> 3 segments (1 local log segment + 2 segments in the remote storage) + // topicRLS4 -> 4 segments (1 local log segment + 3 segments in the remote storage) Map rlsConfigs = new HashMap<>(); rlsConfigs.put(TopicConfig.REMOTE_LOG_STORAGE_ENABLE_CONFIG, "true"); rlsConfigs.put(TopicConfig.LOCAL_LOG_RETENTION_BYTES_CONFIG, "1"); @@ -134,23 +139,22 @@ private static List withRemoteStorage() { serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_STORAGE_SYSTEM_ENABLE_PROP, "true"); serverProperties.put(RemoteLogManagerConfig.REMOTE_STORAGE_MANAGER_CLASS_NAME_PROP, LocalTieredStorage.class.getName()); serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_MANAGER_TASK_INTERVAL_MS_PROP, "1000"); + serverProperties.put(TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_PARTITIONS_PROP, "1"); serverProperties.put(ServerLogConfigs.LOG_CLEANUP_INTERVAL_MS_CONFIG, "1000"); serverProperties.put(ServerLogConfigs.LOG_INITIAL_TASK_DELAY_MS_CONFIG, "100"); - - Map zkProperties = new HashMap<>(serverProperties); - zkProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "PLAINTEXT"); - - Map raftProperties = new HashMap<>(serverProperties); - raftProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "EXTERNAL"); + serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "EXTERNAL"); return Arrays.asList( ClusterConfig.defaultBuilder() .setTypes(Collections.singleton(ZK)) - .setServerProperties(zkProperties) + .setServerProperties(serverProperties) + // align listener name since in KafkaClusterTestKit the default broker listener name is EXTERNAL + // while ZK is PLAINTEXT + .setListenerName("EXTERNAL") .build(), ClusterConfig.defaultBuilder() .setTypes(Stream.of(KRAFT, CO_KRAFT).collect(Collectors.toSet())) - .setServerProperties(raftProperties) + .setServerProperties(serverProperties) .build()); } @@ -362,6 +366,7 @@ public void testGetOffsetsByLatestTieredSpec() throws InterruptedException { executeAndParse("--topic-partitions", "topic\\d+:0", "--time", time)); // test topics enable remote log storage + // topicRLS1 has no result because there's no log segments being uploaded to the remote storage TestUtils.waitForCondition(() -> Arrays.asList( new Row("topicRLS2", 0, 0L), From 7c5d7cc817cee4a3a7988633d3ec03e49ac0682b Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sun, 4 Aug 2024 01:16:38 +0800 Subject: [PATCH 06/10] Add more comment --- .../test/java/org/apache/kafka/tools/GetOffsetShellTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index 831baa8ab7224..f655a1d2f9472 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -148,8 +148,9 @@ private static List withRemoteStorage() { ClusterConfig.defaultBuilder() .setTypes(Collections.singleton(ZK)) .setServerProperties(serverProperties) - // align listener name since in KafkaClusterTestKit the default broker listener name is EXTERNAL - // while ZK is PLAINTEXT + // we set REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP to EXTERNAL, so we need to + // align listener name here as KafkaClusterTestKit (RAFT/CO_RAFT) the default + // broker listener name is EXTERNAL while in ZK it is PLAINTEXT .setListenerName("EXTERNAL") .build(), ClusterConfig.defaultBuilder() From b6426e51d40e5b921f7d497bb0f1c3b2b88adf68 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sun, 4 Aug 2024 01:20:20 +0800 Subject: [PATCH 07/10] Modify comment --- .../test/java/org/apache/kafka/tools/GetOffsetShellTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index f655a1d2f9472..35cacab98da76 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -149,7 +149,7 @@ private static List withRemoteStorage() { .setTypes(Collections.singleton(ZK)) .setServerProperties(serverProperties) // we set REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP to EXTERNAL, so we need to - // align listener name here as KafkaClusterTestKit (RAFT/CO_RAFT) the default + // align listener name here as KafkaClusterTestKit (KRAFT/CO_KRAFT) the default // broker listener name is EXTERNAL while in ZK it is PLAINTEXT .setListenerName("EXTERNAL") .build(), From 52f09704c231ed48e63ad1bc1357438bbc14b2ac Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sun, 4 Aug 2024 11:56:40 +0800 Subject: [PATCH 08/10] Address comments --- .../kafka/clients/admin/KafkaAdminClient.java | 2 +- .../kafka/clients/admin/OffsetSpec.java | 8 ++--- .../clients/admin/KafkaAdminClientTest.java | 4 +-- .../apache/kafka/tools/GetOffsetShell.java | 10 +++---- .../kafka/tools/GetOffsetShellTest.java | 29 ++++++++++++------- 5 files changed, 31 insertions(+), 22 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 8eb7fb4e8c05e..2f195489adda8 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -4862,7 +4862,7 @@ private static long getOffsetFromSpec(OffsetSpec offsetSpec) { return ListOffsetsRequest.MAX_TIMESTAMP; } else if (offsetSpec instanceof OffsetSpec.EarliestLocalSpec) { return ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP; - } else if (offsetSpec instanceof OffsetSpec.LatestTierSpec) { + } else if (offsetSpec instanceof OffsetSpec.LatestTieredSpec) { return ListOffsetsRequest.LATEST_TIERED_TIMESTAMP; } return ListOffsetsRequest.LATEST_TIMESTAMP; diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java b/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java index f2f8f6604022a..68f94cc493e5a 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/OffsetSpec.java @@ -27,7 +27,7 @@ public static class EarliestSpec extends OffsetSpec { } public static class LatestSpec extends OffsetSpec { } public static class MaxTimestampSpec extends OffsetSpec { } public static class EarliestLocalSpec extends OffsetSpec { } - public static class LatestTierSpec extends OffsetSpec { } + public static class LatestTieredSpec extends OffsetSpec { } public static class TimestampSpec extends OffsetSpec { private final long timestamp; @@ -79,7 +79,7 @@ public static OffsetSpec maxTimestamp() { *
* Note: When tiered Storage is not enabled, it behaves the same as retrieving the earliest timestamp offset. */ - public static OffsetSpec earliestLocalSpec() { + public static OffsetSpec earliestLocal() { return new EarliestLocalSpec(); } @@ -88,7 +88,7 @@ public static OffsetSpec earliestLocalSpec() { *
* Note: When tiered storage is not enabled, we will return unknown offset. */ - public static OffsetSpec latestTierSpec() { - return new LatestTierSpec(); + public static OffsetSpec latestTiered() { + return new LatestTieredSpec(); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 8d70e60fc0592..dc74a17844266 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -5864,7 +5864,7 @@ public void testListOffsetsEarliestLocalSpecMinVersion() throws Exception { env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); - env.adminClient().listOffsets(Collections.singletonMap(tp0, OffsetSpec.earliestLocalSpec())); + env.adminClient().listOffsets(Collections.singletonMap(tp0, OffsetSpec.earliestLocal())); TestUtils.waitForCondition(() -> env.kafkaClient().requests().stream().anyMatch(request -> request.requestBuilder().apiKey().messageType == ApiMessageType.LIST_OFFSETS && request.requestBuilder().oldestAllowedVersion() == 9 @@ -5892,7 +5892,7 @@ public void testListOffsetsLatestTierSpecSpecMinVersion() throws Exception { env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); - env.adminClient().listOffsets(Collections.singletonMap(tp0, OffsetSpec.latestTierSpec())); + env.adminClient().listOffsets(Collections.singletonMap(tp0, OffsetSpec.latestTiered())); TestUtils.waitForCondition(() -> env.kafkaClient().requests().stream().anyMatch(request -> request.requestBuilder().apiKey().messageType == ApiMessageType.LIST_OFFSETS && request.requestBuilder().oldestAllowedVersion() == 9 diff --git a/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java b/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java index 9b1ebe9cfe90d..e9e65ba9bf4fc 100644 --- a/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java +++ b/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java @@ -132,7 +132,7 @@ public GetOffsetShellOptions(String[] args) throws TerseException { .ofType(String.class); timeOpt = parser.accepts("time", "timestamp of the offsets before that. [Note: No offset is returned, if the timestamp greater than recently committed record timestamp is given.]") .withRequiredArg() - .describedAs(" / -1 or latest / -2 or earliest / -3 or max-timestamp") + .describedAs(" / -1 or latest / -2 or earliest / -3 or max-timestamp / -4 or earliest-local / -5 or latest-tiered") .ofType(String.class) .defaultsTo("latest"); commandConfigOpt = parser.accepts("command-config", "Property file containing configs to be passed to Admin Client.") @@ -284,9 +284,9 @@ private OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseExce case "max-timestamp": return OffsetSpec.maxTimestamp(); case "earliest-local": - return OffsetSpec.earliestLocalSpec(); + return OffsetSpec.earliestLocal(); case "latest-tiered": - return OffsetSpec.latestTierSpec(); + return OffsetSpec.latestTiered(); default: long timestamp; @@ -304,9 +304,9 @@ private OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseExce } else if (timestamp == ListOffsetsRequest.MAX_TIMESTAMP) { return OffsetSpec.maxTimestamp(); } else if (timestamp == ListOffsetsRequest.EARLIEST_LOCAL_TIMESTAMP) { - return OffsetSpec.earliestLocalSpec(); + return OffsetSpec.earliestLocal(); } else if (timestamp == ListOffsetsRequest.LATEST_TIERED_TIMESTAMP) { - return OffsetSpec.latestTierSpec(); + return OffsetSpec.latestTiered(); } else { return OffsetSpec.forTimestamp(timestamp); } diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index 35cacab98da76..eaea1f0e0f7bf 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -144,18 +144,14 @@ private static List withRemoteStorage() { serverProperties.put(ServerLogConfigs.LOG_INITIAL_TASK_DELAY_MS_CONFIG, "100"); serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "EXTERNAL"); - return Arrays.asList( + return Collections.singletonList( + // we set REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP to EXTERNAL, so we need to + // align listener name here as KafkaClusterTestKit (KRAFT/CO_KRAFT) the default + // broker listener name is EXTERNAL while in ZK it is PLAINTEXT ClusterConfig.defaultBuilder() - .setTypes(Collections.singleton(ZK)) + .setTypes(Stream.of(ZK, KRAFT, CO_KRAFT).collect(Collectors.toSet())) .setServerProperties(serverProperties) - // we set REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP to EXTERNAL, so we need to - // align listener name here as KafkaClusterTestKit (KRAFT/CO_KRAFT) the default - // broker listener name is EXTERNAL while in ZK it is PLAINTEXT .setListenerName("EXTERNAL") - .build(), - ClusterConfig.defaultBuilder() - .setTypes(Stream.of(KRAFT, CO_KRAFT).collect(Collectors.toSet())) - .setServerProperties(serverProperties) .build()); } @@ -340,9 +336,22 @@ public void testGetOffsetsByMaxTimestamp() { @ClusterTemplate("withRemoteStorage") public void testGetOffsetsByEarliestLocalSpec() throws InterruptedException { + setUp(); setUpRemoteLogTopics(); for (String time : new String[] {"-4", "earliest-local"}) { + // test topics disable remote log storage + // as remote log disabled, broker return the same result as earliest offset + TestUtils.waitForCondition(() -> + Arrays.asList( + new Row("topic1", 0, 0L), + new Row("topic2", 0, 0L), + new Row("topic3", 0, 0L), + new Row("topic4", 0, 0L)) + .equals(executeAndParse("--topic-partitions", "topic\\d+.*:0", "--time", time)), + "testGetOffsetsByEarliestLocalSpec get topics with remote log disabled result not match"); + + // test topics enable remote log storage TestUtils.waitForCondition(() -> Arrays.asList( new Row("topicRLS1", 0, 0L), @@ -350,7 +359,7 @@ public void testGetOffsetsByEarliestLocalSpec() throws InterruptedException { new Row("topicRLS3", 0, 2L), new Row("topicRLS4", 0, 3L)) .equals(executeAndParse("--topic-partitions", "topicRLS.*:0", "--time", time)), - "testGetOffsetsByEarliestLocalSpec result not match"); + "testGetOffsetsByEarliestLocalSpec get topics with remote log enabled result not match"); } } From d23ef673f63cdc11442bfe6b61e060e8da7197a5 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sun, 4 Aug 2024 13:02:15 +0800 Subject: [PATCH 09/10] Address comment --- .../main/java/org/apache/kafka/tools/GetOffsetShell.java | 5 +++-- .../org/apache/kafka/tools/GetOffsetShellParsingTest.java | 8 ++++++++ 2 files changed, 11 insertions(+), 2 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java b/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java index e9e65ba9bf4fc..60b78acd22be4 100644 --- a/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java +++ b/tools/src/main/java/org/apache/kafka/tools/GetOffsetShell.java @@ -275,7 +275,8 @@ public Map fetchOffsets(GetOffsetShellOptions options) thr } } - private OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseException { + // visible for tseting + static OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseException { switch (listOffsetsTimestamp) { case "earliest": return OffsetSpec.earliest(); @@ -294,7 +295,7 @@ private OffsetSpec parseOffsetSpec(String listOffsetsTimestamp) throws TerseExce timestamp = Long.parseLong(listOffsetsTimestamp); } catch (NumberFormatException e) { throw new TerseException("Malformed time argument " + listOffsetsTimestamp + ". " + - "Please use -1 or latest / -2 or earliest / -3 or max-timestamp, or a specified long format timestamp"); + "Please use -1 or latest / -2 or earliest / -3 or max-timestamp / -4 or earliest-local / -5 or latest-tiered, or a specified long format timestamp"); } if (timestamp == ListOffsetsRequest.EARLIEST_TIMESTAMP) { diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellParsingTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellParsingTest.java index 3c4ef0894f71c..9e81c23f3099e 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellParsingTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellParsingTest.java @@ -22,6 +22,7 @@ import org.junit.jupiter.api.Test; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -243,6 +244,13 @@ public void testInvalidTimeValue() { assertThrows(TerseException.class, () -> GetOffsetShell.execute("--bootstrap-server", "localhost:9092", "--time", "invalid")); } + @Test + public void testInvalidOffset() { + assertEquals("Malformed time argument foo. " + + "Please use -1 or latest / -2 or earliest / -3 or max-timestamp / -4 or earliest-local / -5 or latest-tiered, or a specified long format timestamp", + assertThrows(TerseException.class, () -> GetOffsetShell.parseOffsetSpec("foo")).getMessage()); + } + private TopicPartition getTopicPartition(String topic, Integer partition) { return new TopicPartition(topic, partition); } From ce1e59d9763d03491a9709b3a8ae9653ad1e69a6 Mon Sep 17 00:00:00 2001 From: Kuan-Po Tseng Date: Sun, 4 Aug 2024 16:26:47 +0800 Subject: [PATCH 10/10] Add missing config prefix --- .../test/java/org/apache/kafka/tools/GetOffsetShellTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java index eaea1f0e0f7bf..95007d7bf8541 100644 --- a/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/GetOffsetShellTest.java @@ -136,10 +136,10 @@ private void sendProducerRecords(Function topicName) { private static List withRemoteStorage() { Map serverProperties = new HashMap<>(); serverProperties.put(RemoteLogManagerConfig.DEFAULT_REMOTE_LOG_METADATA_MANAGER_CONFIG_PREFIX + TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP, "1"); + serverProperties.put(RemoteLogManagerConfig.DEFAULT_REMOTE_LOG_METADATA_MANAGER_CONFIG_PREFIX + TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_PARTITIONS_PROP, "1"); serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_STORAGE_SYSTEM_ENABLE_PROP, "true"); serverProperties.put(RemoteLogManagerConfig.REMOTE_STORAGE_MANAGER_CLASS_NAME_PROP, LocalTieredStorage.class.getName()); serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_MANAGER_TASK_INTERVAL_MS_PROP, "1000"); - serverProperties.put(TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_PARTITIONS_PROP, "1"); serverProperties.put(ServerLogConfigs.LOG_CLEANUP_INTERVAL_MS_CONFIG, "1000"); serverProperties.put(ServerLogConfigs.LOG_INITIAL_TASK_DELAY_MS_CONFIG, "100"); serverProperties.put(RemoteLogManagerConfig.REMOTE_LOG_METADATA_MANAGER_LISTENER_NAME_PROP, "EXTERNAL");