Skip to content

Commit 1c473c1

Browse files
committed
Use ClusterInstance#createTopic(WithAssignment) in tools tests
1 parent 596abe3 commit 1c473c1

4 files changed

Lines changed: 30 additions & 45 deletions

File tree

‎tools/src/test/java/org/apache/kafka/tools/ConfigCommandIntegrationTest.java‎

Lines changed: 13 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -217,23 +217,19 @@ public void testNullStatusOnKraftCommandAlterClientMetrics() {
217217

218218
@ClusterTest
219219
public void testAddConfigKeyValuesUsingCommand() throws Exception {
220-
try (Admin client = cluster.admin()) {
221-
NewTopic newTopic = new NewTopic("topic", 1, (short) 1);
222-
client.createTopics(Set.of(newTopic)).all().get();
223-
cluster.waitTopicCreation("topic", 1);
224-
Stream<String> command = Stream.concat(quorumArgs(), Stream.of(
225-
"--entity-type", "topics",
226-
"--entity-name", "topic",
227-
"--alter", "--add-config", "cleanup.policy=[delete,compact]"));
228-
String message = captureStandardOut(run(command));
229-
assertEquals("Completed updating config for topic topic.", message);
230-
command = Stream.concat(quorumArgs(), Stream.of(
231-
"--entity-type", "topics",
232-
"--entity-name", "topic",
233-
"--describe"));
234-
message = captureStandardOut(run(command));
235-
assertTrue(message.contains("cleanup.policy=delete,compact"), "Config entry was not added correctly");
236-
}
220+
cluster.createTopic("topic", 1, (short) 1);
221+
Stream<String> command = Stream.concat(quorumArgs(), Stream.of(
222+
"--entity-type", "topics",
223+
"--entity-name", "topic",
224+
"--alter", "--add-config", "cleanup.policy=[delete,compact]"));
225+
String message = captureStandardOut(run(command));
226+
assertEquals("Completed updating config for topic topic.", message);
227+
command = Stream.concat(quorumArgs(), Stream.of(
228+
"--entity-type", "topics",
229+
"--entity-name", "topic",
230+
"--describe"));
231+
message = captureStandardOut(run(command));
232+
assertTrue(message.contains("cleanup.policy=delete,compact"), "Config entry was not added correctly");
237233
}
238234

239235
@ClusterTest

‎tools/src/test/java/org/apache/kafka/tools/LogDirsCommandTest.java‎

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,6 @@
2020
import org.apache.kafka.clients.admin.Admin;
2121
import org.apache.kafka.clients.admin.LogDirDescription;
2222
import org.apache.kafka.clients.admin.MockAdminClient;
23-
import org.apache.kafka.clients.admin.NewTopic;
2423
import org.apache.kafka.common.Node;
2524
import org.apache.kafka.common.TopicPartition;
2625
import org.apache.kafka.common.test.ClusterInstance;
@@ -218,9 +217,6 @@ private String execute(LogDirsCommand.LogDirsCommandOptions options, Admin admin
218217
}
219218

220219
private void createTopic(ClusterInstance clusterInstance, String topic) {
221-
try (Admin admin = Admin.create(Map.of(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers()))) {
222-
assertDoesNotThrow(() -> admin.createTopics(List.of(new NewTopic(topic, Map.of(0, List.of(0))))).topicId(topic).get());
223-
assertDoesNotThrow(() -> clusterInstance.waitTopicCreation(topic, 1));
224-
}
220+
assertDoesNotThrow(() -> clusterInstance.createTopicWithAssignment(topic, Map.of(0, List.of(0))));
225221
}
226222
}

‎tools/src/test/java/org/apache/kafka/tools/consumer/group/ResetConsumerGroupOffsetTest.java‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -160,10 +160,9 @@ public void testResetOffsetsWithOfflinePartitionNotInResetTarget(ClusterInstance
160160
String group = "new.group";
161161
String[] args = buildArgsForGroup(cluster, group, "--to-earliest", "--execute", "--topic", topic + ":0");
162162

163-
try (Admin admin = cluster.admin(); ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) {
164-
admin.createTopics(List.of(new NewTopic(topic, Map.of(0, List.of(0), 1, List.of(1)))));
165-
cluster.waitTopicCreation(topic, 2);
163+
cluster.createTopicWithAssignment(topic, Map.of(0, List.of(0), 1, List.of(1)));
166164

165+
try (ConsumerGroupCommand.ConsumerGroupService service = getConsumerGroupService(args)) {
167166
cluster.shutdownBroker(1);
168167

169168
Map<TopicPartition, OffsetAndMetadata> resetOffsets = service.resetOffsets().get(group);

‎tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java‎

Lines changed: 14 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@
2424
import org.apache.kafka.clients.admin.ConfigEntry;
2525
import org.apache.kafka.clients.admin.DescribeLogDirsResult;
2626
import org.apache.kafka.clients.admin.ListOffsetsResult.ListOffsetsResultInfo;
27-
import org.apache.kafka.clients.admin.NewTopic;
2827
import org.apache.kafka.clients.admin.OffsetSpec;
2928
import org.apache.kafka.clients.admin.TopicDescription;
3029
import org.apache.kafka.clients.consumer.Consumer;
@@ -546,25 +545,20 @@ public void testExecuteAssignmentWithOneBootstrapServerShutdownWontTimeout() thr
546545
}
547546

548547
private void createTopics() {
549-
try (Admin admin = Admin.create(Map.of(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers()))) {
550-
Map<Integer, List<Integer>> fooReplicasAssignments = new HashMap<>();
551-
fooReplicasAssignments.put(0, List.of(0, 1, 2));
552-
fooReplicasAssignments.put(1, List.of(1, 2, 3));
553-
Assertions.assertDoesNotThrow(() -> admin.createTopics(List.of(new NewTopic("foo", fooReplicasAssignments))).topicId("foo").get());
554-
Assertions.assertDoesNotThrow(() -> clusterInstance.waitTopicCreation("foo", fooReplicasAssignments.size()));
555-
556-
Map<Integer, List<Integer>> barReplicasAssignments = new HashMap<>();
557-
barReplicasAssignments.put(0, List.of(3, 2, 1));
558-
Assertions.assertDoesNotThrow(() -> admin.createTopics(List.of(new NewTopic("bar", barReplicasAssignments))).topicId("bar").get());
559-
Assertions.assertDoesNotThrow(() -> clusterInstance.waitTopicCreation("bar", barReplicasAssignments.size()));
560-
561-
Map<Integer, List<Integer>> bazReplicasAssignments = new HashMap<>();
562-
bazReplicasAssignments.put(0, List.of(1, 0, 2));
563-
bazReplicasAssignments.put(1, List.of(2, 0, 1));
564-
bazReplicasAssignments.put(2, List.of(0, 2, 1));
565-
Assertions.assertDoesNotThrow(() -> admin.createTopics(List.of(new NewTopic("baz", bazReplicasAssignments))).topicId("baz").get());
566-
Assertions.assertDoesNotThrow(() -> clusterInstance.waitTopicCreation("baz", bazReplicasAssignments.size()));
567-
}
548+
Map<Integer, List<Integer>> fooReplicasAssignments = new HashMap<>();
549+
fooReplicasAssignments.put(0, List.of(0, 1, 2));
550+
fooReplicasAssignments.put(1, List.of(1, 2, 3));
551+
Assertions.assertDoesNotThrow(() -> clusterInstance.createTopicWithAssignment("foo", fooReplicasAssignments));
552+
553+
Map<Integer, List<Integer>> barReplicasAssignments = new HashMap<>();
554+
barReplicasAssignments.put(0, List.of(3, 2, 1));
555+
Assertions.assertDoesNotThrow(() -> clusterInstance.createTopicWithAssignment("bar", barReplicasAssignments));
556+
557+
Map<Integer, List<Integer>> bazReplicasAssignments = new HashMap<>();
558+
bazReplicasAssignments.put(0, List.of(1, 0, 2));
559+
bazReplicasAssignments.put(1, List.of(2, 0, 1));
560+
bazReplicasAssignments.put(2, List.of(0, 2, 1));
561+
Assertions.assertDoesNotThrow(() -> clusterInstance.createTopicWithAssignment("baz", bazReplicasAssignments));
568562
}
569563

570564
private void produceMessages(String topic, int partition, int numMessages) {

0 commit comments

Comments
 (0)