Skip to content

Commit 596abe3

Browse files
committed
Use ClusterInstance#createTopic(WithAssignment) in remote log metadata manager tests
1 parent d61de01 commit 596abe3

3 files changed

Lines changed: 10 additions & 29 deletions

File tree

‎storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerMultipleSubscriptionsTest.java‎

Lines changed: 1 addition & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,6 @@
1717
package org.apache.kafka.server.log.remote.metadata.storage;
1818

1919

20-
import org.apache.kafka.clients.CommonClientConfigs;
21-
import org.apache.kafka.clients.admin.Admin;
22-
import org.apache.kafka.clients.admin.NewTopic;
2320
import org.apache.kafka.common.TopicIdPartition;
2421
import org.apache.kafka.common.TopicPartition;
2522
import org.apache.kafka.common.Uuid;
@@ -171,9 +168,6 @@ public int metadataPartition(TopicIdPartition topicIdPartition) {
171168
}
172169

173170
private void createTopic(String topic, Map<Integer, List<Integer>> replicasAssignments) {
174-
try (Admin admin = Admin.create(Map.of(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers()))) {
175-
admin.createTopics(List.of(new NewTopic(topic, replicasAssignments)));
176-
assertDoesNotThrow(() -> clusterInstance.waitTopicCreation(topic, replicasAssignments.size()));
177-
}
171+
assertDoesNotThrow(() -> clusterInstance.createTopicWithAssignment(topic, replicasAssignments));
178172
}
179173
}

‎storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerRestartTest.java‎

Lines changed: 4 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,6 @@
1616
*/
1717
package org.apache.kafka.server.log.remote.metadata.storage;
1818

19-
import org.apache.kafka.clients.admin.Admin;
20-
import org.apache.kafka.clients.admin.NewTopic;
2119
import org.apache.kafka.common.TopicIdPartition;
2220
import org.apache.kafka.common.TopicPartition;
2321
import org.apache.kafka.common.Uuid;
@@ -58,15 +56,10 @@ public void testRLMMAPIsAfterRestart() throws Exception {
5856
// Create topics.
5957
String leaderTopic = "new-leader";
6058
String followerTopic = "new-follower";
61-
try (Admin admin = clusterInstance.admin()) {
62-
// Set broker id 0 as the first entry which is taken as the leader.
63-
NewTopic newLeaderTopic = new NewTopic(leaderTopic, Map.of(0, List.of(0, 1, 2)));
64-
// Set broker id 1 as the first entry which is taken as the leader.
65-
NewTopic newFollowerTopic = new NewTopic(followerTopic, Map.of(0, List.of(1, 2, 0)));
66-
admin.createTopics(List.of(newLeaderTopic, newFollowerTopic)).all().get();
67-
}
68-
clusterInstance.waitTopicCreation(leaderTopic, 1);
69-
clusterInstance.waitTopicCreation(followerTopic, 1);
59+
// Set broker id 0 as the first entry which is taken as the leader.
60+
clusterInstance.createTopicWithAssignment(leaderTopic, Map.of(0, List.of(0, 1, 2)));
61+
// Set broker id 1 as the first entry which is taken as the leader.
62+
clusterInstance.createTopicWithAssignment(followerTopic, Map.of(0, List.of(1, 2, 0)));
7063

7164
TopicIdPartition leaderTopicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition(leaderTopic, 0));
7265
TopicIdPartition followerTopicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition(followerTopic, 0));

‎storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java‎

Lines changed: 5 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121
import org.apache.kafka.clients.admin.ConfigEntry;
2222
import org.apache.kafka.clients.admin.DescribeConfigsResult;
2323
import org.apache.kafka.clients.admin.DescribeTopicsResult;
24-
import org.apache.kafka.clients.admin.NewTopic;
2524
import org.apache.kafka.clients.admin.TopicDescription;
2625
import org.apache.kafka.common.KafkaFuture;
2726
import org.apache.kafka.common.TopicIdPartition;
@@ -94,10 +93,9 @@ public void teardown() throws IOException {
9493

9594
@ClusterTest
9695
public void testDoesTopicExist() throws ExecutionException, InterruptedException {
96+
String topic = "test-topic-exist";
97+
clusterInstance.createTopic(topic, 1, (short) 1);
9798
try (Admin admin = clusterInstance.admin()) {
98-
String topic = "test-topic-exist";
99-
admin.createTopics(List.of(new NewTopic(topic, 1, (short) 1))).all().get();
100-
clusterInstance.waitTopicCreation(topic, 1);
10199
boolean doesTopicExist = topicBasedRlmm().doesTopicExist(admin, topic);
102100
assertTrue(doesTopicExist);
103101
}
@@ -144,13 +142,9 @@ public void testNewPartitionUpdates() throws Exception {
144142
// Create topics.
145143
String leaderTopic = "new-leader";
146144
String followerTopic = "new-follower";
147-
try (Admin admin = clusterInstance.admin()) {
148-
// Set broker id 0 as the first entry which is taken as the leader.
149-
admin.createTopics(List.of(new NewTopic(leaderTopic, Map.of(0, List.of(0, 1, 2))))).all().get();
150-
clusterInstance.waitTopicCreation(leaderTopic, 1);
151-
admin.createTopics(List.of(new NewTopic(followerTopic, Map.of(0, List.of(1, 2, 0))))).all().get();
152-
clusterInstance.waitTopicCreation(followerTopic, 1);
153-
}
145+
// Set broker id 0 as the first entry which is taken as the leader.
146+
clusterInstance.createTopicWithAssignment(leaderTopic, Map.of(0, List.of(0, 1, 2)));
147+
clusterInstance.createTopicWithAssignment(followerTopic, Map.of(0, List.of(1, 2, 0)));
154148

155149
final TopicIdPartition newLeaderTopicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition(leaderTopic, 0));
156150
final TopicIdPartition newFollowerTopicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition(followerTopic, 0));

0 commit comments

Comments
 (0)