Skip to content
Merged
Show file tree
Hide file tree
Changes from 12 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -119,9 +119,7 @@ public void testAdminRebootstrapDisabled(ClusterInstance clusterInstance) throws
}
)
public void testProducerRebootstrap(ClusterInstance clusterInstance) throws ExecutionException, InterruptedException {
try (var admin = clusterInstance.admin()) {
admin.createTopics(List.of(new NewTopic(TOPIC, PARTITIONS, (short) REPLICAS)));
}
clusterInstance.createTopic(TOPIC, PARTITIONS, (short) REPLICAS);

var broker0 = 0;
var broker1 = 1;
Expand Down Expand Up @@ -154,9 +152,7 @@ public void testProducerRebootstrap(ClusterInstance clusterInstance) throws Exec
}
)
public void testProducerRebootstrapDisabled(ClusterInstance clusterInstance) throws ExecutionException, InterruptedException {
try (var admin = clusterInstance.admin()) {
admin.createTopics(List.of(new NewTopic(TOPIC, PARTITIONS, (short) REPLICAS)));
}
clusterInstance.createTopic(TOPIC, PARTITIONS, (short) REPLICAS);

var broker0 = 0;
var broker1 = 1;
Expand Down Expand Up @@ -307,9 +303,7 @@ public void testConsumerRebootstrapDisabled(ClusterInstance clusterInstance) thr
}
)
public void testRebootstrapOnMetadataClusterCheckFail(ClusterInstance clusterInstance) throws ExecutionException, InterruptedException {
try (var admin = clusterInstance.admin()) {
admin.createTopics(List.of(new NewTopic(TOPIC, 2, (short) 1)));
}
clusterInstance.createTopic(TOPIC, 2, (short) 1);

try (var producer = clusterInstance.producer()) {
var recordMetadata0 = producer.send(new ProducerRecord<>(TOPIC, 0, null, "value 0".getBytes())).get();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -130,14 +130,10 @@ public void testIncrementPartitions(ClusterInstance cluster) throws Exception {
@ClusterTest(brokers = 3, controllers = 3, metadataVersion = MetadataVersion.IBP_3_7_IV2)
})
public void testCreatePartitionsAcrossMetadataVersions(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
Map<String, KafkaFuture<Void>> createResults = admin.createTopics(List.of(
new NewTopic("foo", 1, (short) 3),
new NewTopic("bar", 2, (short) 3)
)).values();
createResults.get("foo").get();
createResults.get("bar").get();
cluster.createTopic("foo", 1, (short) 3);
cluster.createTopic("bar", 2, (short) 3);

try (Admin admin = cluster.admin()) {
Map<String, KafkaFuture<Void>> increaseResults = admin.createPartitions(Map.of(
"foo", NewPartitions.increaseTo(3),
"bar", NewPartitions.increaseTo(2)
Expand Down Expand Up @@ -246,11 +242,8 @@ public void testCreatePartitions(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
String topic1 = "create-partitions-topic-1";
String topic2 = "create-partitions-topic-2";
admin.createTopics(List.of(
new NewTopic(topic1, 1, (short) 1),
new NewTopic(topic2, 1, (short) 2))).all().get();
cluster.waitTopicCreation(topic1, 1);
cluster.waitTopicCreation(topic2, 1);
cluster.createTopic(topic1, 1, (short) 1);
cluster.createTopic(topic2, 1, (short) 2);
assertEquals(1, numPartitions(admin, topic1));
assertEquals(1, numPartitions(admin, topic2));

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -277,11 +277,9 @@ public void testDescribeReplicaLogDirs() throws Exception {

@ClusterTest
public void testDescribeTopicsWithOptionPartitionSizeLimitPerResponse() throws Exception {
String testTopic = "test-topic";
clusterInstance.createTopic(testTopic, 3, (short) 1);
try (Admin admin = clusterInstance.admin()) {
String testTopic = "test-topic";
admin.createTopics(List.of(new NewTopic(testTopic, 3, (short) 1))).all().get();
clusterInstance.waitTopicCreation(testTopic, 3);

Map<String, TopicDescription> topics = admin.describeTopics(List.of(testTopic),
new DescribeTopicsOptions().partitionSizeLimitPerResponse(1)).allTopicNames().get();
assertEquals(1, topics.size());
Expand All @@ -307,11 +305,9 @@ public void testDescribeTopicsWithOptionTimeoutMs() {
*/
@ClusterTest
public void testDescribeNonExistingTopic() throws Exception {
String existingTopic = "existing-topic";
clusterInstance.createTopic(existingTopic, 1, (short) 1);
try (Admin admin = clusterInstance.admin()) {
String existingTopic = "existing-topic";
admin.createTopics(List.of(new NewTopic(existingTopic, 1, (short) 1))).all().get();
clusterInstance.waitTopicCreation(existingTopic, 1);

String nonExistingTopic = "non-existing";
Map<String, KafkaFuture<TopicDescription>> results =
admin.describeTopics(List.of(nonExistingTopic, existingTopic)).topicNameValues();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,11 +70,9 @@ public void testClientInstanceId(ClusterInstance clusterInstance) throws Interru
Map<String, Object> configs = new HashMap<>();
configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers());
configs.put(AdminClientConfig.ENABLE_METRICS_PUSH_CONFIG, true);
String testTopicName = "test_topic";
clusterInstance.createTopic(testTopicName, 1, (short) 1);
try (Admin admin = Admin.create(configs)) {
String testTopicName = "test_topic";
admin.createTopics(Collections.singletonList(new NewTopic(testTopicName, 1, (short) 1)));
clusterInstance.waitTopicCreation(testTopicName, 1);

Map<String, Object> producerConfigs = new HashMap<>();
producerConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers());
producerConfigs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,10 +71,7 @@ public void testCreateAndDeleteTopic(ClusterInstance cluster) throws Exception {
String testTopic = "test-topic";
try (Admin admin = cluster.admin()) {
// Create a test topic
List<NewTopic> newTopics = List.of(new NewTopic(testTopic, 1, (short) 3));
CreateTopicsResult createTopicResult = admin.createTopics(newTopics);
createTopicResult.all().get();
cluster.waitTopicCreation(testTopic, 1);
cluster.createTopic(testTopic, 1, (short) 3);

// Delete topic
DeleteTopicsResult deleteResult = admin.deleteTopics(List.of(testTopic));
Expand All @@ -88,7 +85,7 @@ public void testCreateAndDeleteTopic(ClusterInstance cluster) throws Exception {
@ClusterTest
public void testDeleteTopicWithAllAliveReplicas(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
admin.deleteTopics(List.of(DEFAULT_TOPIC)).all().get();
cluster.waitTopicDeletion(DEFAULT_TOPIC);
}
Expand All @@ -97,7 +94,7 @@ public void testDeleteTopicWithAllAliveReplicas(ClusterInstance cluster) throws
@ClusterTest
public void testResumeDeleteTopicWithRecoveredFollower(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
int leaderId = waitUtilLeaderIsKnown(cluster.brokers(), topicPartition);
KafkaBroker follower = findFollower(cluster.brokers().values(), leaderId);
Expand All @@ -120,7 +117,7 @@ public void testResumeDeleteTopicWithRecoveredFollower(ClusterInstance cluster)
@ClusterTest(brokers = 4)
public void testPartitionReassignmentDuringDeleteTopic(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
Map<Integer, KafkaBroker> servers = findPartitionHostingBrokers(cluster.brokers());
int leaderId = waitUtilLeaderIsKnown(cluster.brokers(), topicPartition);
Expand All @@ -146,7 +143,7 @@ public void testPartitionReassignmentDuringDeleteTopic(ClusterInstance cluster)
@ClusterTest(brokers = 4)
public void testIncreasePartitionCountDuringDeleteTopic(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
Map<Integer, KafkaBroker> partitionHostingBrokers = findPartitionHostingBrokers(cluster.brokers());
waitForReplicaCreated(partitionHostingBrokers, topicPartition, "Replicas for topic test not created.");
Expand Down Expand Up @@ -174,7 +171,7 @@ public void testIncreasePartitionCountDuringDeleteTopic(ClusterInstance cluster)
@ClusterTest
public void testDeleteTopicDuringAddPartition(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
int leaderId = waitUtilLeaderIsKnown(cluster.brokers(), new TopicPartition(DEFAULT_TOPIC, 0));
TopicPartition newTopicPartition = new TopicPartition(DEFAULT_TOPIC, 1);
KafkaBroker follower = findFollower(cluster.brokers().values(), leaderId);
Expand All @@ -199,7 +196,7 @@ public void testDeleteTopicDuringAddPartition(ClusterInstance cluster) throws Ex
@ClusterTest
public void testAddPartitionDuringDeleteTopic(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
// partitions to be added to the topic later
TopicPartition newTopicPartition = new TopicPartition(DEFAULT_TOPIC, 1);
admin.deleteTopics(List.of(DEFAULT_TOPIC)).all().get();
Expand All @@ -213,20 +210,20 @@ public void testAddPartitionDuringDeleteTopic(ClusterInstance cluster) throws Ex
@ClusterTest
public void testRecreateTopicAfterDeletion(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
admin.deleteTopics(List.of(DEFAULT_TOPIC)).all().get();
cluster.waitTopicDeletion(DEFAULT_TOPIC);
// re-create topic on same replicas
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
waitForReplicaCreated(cluster.brokers(), topicPartition, "Replicas for topic " + DEFAULT_TOPIC + " not created.");
}
}

@ClusterTest
public void testDeleteNonExistingTopic(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
String topic = "test2";
TestUtils.waitForCondition(() -> {
Expand All @@ -252,7 +249,7 @@ public void testDeleteNonExistingTopic(ClusterInstance cluster) throws Exception
})
public void testDeleteTopicWithCleaner(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
// for simplicity, we are validating cleaner offsets on a single broker
KafkaBroker server = cluster.brokers().values().stream().findFirst().orElseThrow();
Expand All @@ -273,7 +270,7 @@ public void testDeleteTopicWithCleaner(ClusterInstance cluster) throws Exception
@ClusterTest
public void testDeleteTopicAlreadyMarkedAsDeleted(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
admin.deleteTopics(List.of(DEFAULT_TOPIC)).all().get();

TestUtils.waitForCondition(() -> {
Expand All @@ -293,7 +290,7 @@ public void testDeleteTopicAlreadyMarkedAsDeleted(ClusterInstance cluster) throw
serverProperties = {@ClusterConfigProperty(key = ServerConfigs.DELETE_TOPIC_ENABLE_CONFIG, value = "false")})
public void testDisableDeleteTopic(ClusterInstance cluster) throws Exception {
try (Admin admin = cluster.admin()) {
admin.createTopics(List.of(new NewTopic(DEFAULT_TOPIC, expectedReplicaAssignment))).all().get();
cluster.createTopicWithAssignment(DEFAULT_TOPIC, expectedReplicaAssignment);
TopicPartition topicPartition = new TopicPartition(DEFAULT_TOPIC, 0);
TestUtils.waitForCondition(() -> {
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -117,8 +117,7 @@ public void testConsumerGroupAuthorizedOperations(ClusterInstance clusterInstanc
try (Admin admin = clusterInstance.admin(createAdminConfig(JaasUtils.KAFKA_PLAIN_ADMIN, JaasUtils.KAFKA_PLAIN_ADMIN_PASSWORD));
Admin user1 = clusterInstance.admin(createAdminConfig(JaasUtils.KAFKA_PLAIN_USER1, JaasUtils.KAFKA_PLAIN_USER1_PASSWORD))
) {
admin.createTopics(List.of(new NewTopic("topic1", 1, (short) 1)));
clusterInstance.waitTopicCreation("topic1", 1);
clusterInstance.createTopic("topic1", 1, (short) 1);

// create consumers to avoid group not found error
TopicPartition tp = new TopicPartition("topic1", 0);
Expand Down Expand Up @@ -188,14 +187,8 @@ public void testTopicAuthorizedOperations(ClusterInstance clusterInstance) throw
String topic1 = "topic1";
String topic2 = "topic2";
setupSecurity(clusterInstance);
try (Admin admin = clusterInstance.admin(createAdminConfig(JaasUtils.KAFKA_PLAIN_ADMIN, JaasUtils.KAFKA_PLAIN_ADMIN_PASSWORD))) {
admin.createTopics(List.of(
new NewTopic(topic1, 1, (short) 1),
new NewTopic(topic2, 1, (short) 1)
));
clusterInstance.waitTopicCreation(topic1, 1);
clusterInstance.waitTopicCreation(topic2, 1);
}
clusterInstance.createTopic(topic1, 1, (short) 1);
clusterInstance.createTopic(topic2, 1, (short) 1);

try (Admin admin = clusterInstance.admin(createAdminConfig(JaasUtils.KAFKA_PLAIN_USER1, JaasUtils.KAFKA_PLAIN_USER1_PASSWORD))) {
// test without includeAuthorizedOperations flag
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
import org.apache.kafka.common.test.api.ClusterTest;
import org.apache.kafka.common.test.api.Type;
import org.apache.kafka.server.config.ServerConfigs;
import org.apache.kafka.test.TestUtils;

import java.util.List;
import java.util.Map;
Expand All @@ -40,11 +39,7 @@ public void testIncreaseNumIoThreads(ClusterInstance cluster) throws Exception {

// KAFKA-16649 introduced this coverage as a safeguard against future deadlocks.
// Verify that the broker continues serving requests after resizing the I/O thread pool.
admin.createTopics(List.of(new NewTopic("test-topic", 1, (short) 1))).all().get();
TestUtils.waitForCondition(
() -> admin.listTopics().names().get().contains("test-topic"),
"Failed to find test-topic"
);
cluster.createTopic("test-topic", 1, (short) 1);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -131,9 +131,7 @@ public void testInternalConfigsDoNotReturnForDescribeConfigs(ClusterInstance clu
ConfigResource groupResource = new ConfigResource(ConfigResource.Type.GROUP, "testGroup");
ConfigResource clientMetricsResource = new ConfigResource(ConfigResource.Type.CLIENT_METRICS, "testClient");

admin.createTopics(List.of(new NewTopic(TOPIC, 1, (short) 1))).config(TOPIC).get();
// make sure the topic metadata exist
cluster.waitTopicCreation(TOPIC, 1);
cluster.createTopic(TOPIC, 1, (short) 1);
Map<ConfigResource, Config> configResourceMap = admin.describeConfigs(
List.of(brokerResource, topicResource, groupResource, clientMetricsResource)).all().get();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@
import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.NewPartitionReassignment;
import org.apache.kafka.clients.admin.NewPartitions;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.internals.AbstractHeartbeatRequestManager;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
Expand Down Expand Up @@ -261,7 +260,9 @@ public void testLeaderEpoch(ClusterInstance clusterInstance) throws Exception {
)
})
public void testRackAwareAssignment(ClusterInstance clusterInstance) throws ExecutionException, InterruptedException {
// Create a new topic with 1 partition on broker 0.
String topic = "test-topic";
clusterInstance.createTopicWithAssignment(topic, Map.of(0, List.of(0)));
try (Admin admin = clusterInstance.admin();
Producer<byte[], byte[]> producer = clusterInstance.producer();
Consumer<byte[], byte[]> consumer0 = clusterInstance.consumer(Map.of(
Expand All @@ -283,10 +284,6 @@ public void testRackAwareAssignment(ClusterInstance clusterInstance) throws Exec
ConsumerConfig.GROUP_PROTOCOL_CONFIG, GroupProtocol.CONSUMER.name()
))
) {
// Create a new topic with 1 partition on broker 0.
admin.createTopics(List.of(new NewTopic(topic, Map.of(0, List.of(0)))));
clusterInstance.waitTopicCreation(topic, 1);

producer.send(new ProducerRecord<>(topic, "key".getBytes(), "value".getBytes()));
producer.flush();

Expand Down Expand Up @@ -376,9 +373,7 @@ public void testSingleCoordinatorOwnershipAfterPartitionReassignment(ClusterInst
producer.send(new ProducerRecord<>("topic", "value".getBytes(StandardCharsets.UTF_8)));
}

try (var admin = clusterInstance.admin()) {
admin.createTopics(List.of(new NewTopic(Topic.GROUP_METADATA_TOPIC_NAME, Map.of(0, List.of(0))))).all().get();
}
clusterInstance.createTopicWithAssignment(Topic.GROUP_METADATA_TOPIC_NAME, Map.of(0, List.of(0)));

try (var consumer = clusterInstance.consumer(Map.of(ConsumerConfig.GROUP_ID_CONFIG, "test-group"));
var admin = clusterInstance.admin()) {
Expand Down
Loading
Loading