This is an automated email from the ASF dual-hosted git repository. MartijnVisser pushed a commit to branch v5.0 in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
commit a3341cc08415c683bb35699f7ac00a1e541b09c9 Author: Aleksandr Savonin <[email protected]> AuthorDate: Mon May 18 19:58:28 2026 +0200 [FLINK-39699][tests] Wait for partitions assignment in KafkaSinkITCase (cherry picked from commit 443c0f3215ae60b69a0049d7e4dad4dcc582f0bc) --- .../flink/connector/kafka/sink/KafkaSinkITCase.java | 21 ++++++++------------- 1 file changed, 8 insertions(+), 13 deletions(-) diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java index 83ed6ac6..19ab1c6a 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java @@ -73,9 +73,7 @@ import org.apache.flink.testutils.junit.SharedReference; import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.clients.admin.AdminClient; -import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.DeleteTopicsResult; -import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.assertj.core.api.InstanceOfAssertFactories; import org.junit.jupiter.api.AfterAll; @@ -125,6 +123,7 @@ import java.util.stream.Stream; import static org.apache.flink.configuration.StateRecoveryOptions.SAVEPOINT_PATH; import static org.apache.flink.connector.kafka.testutils.KafkaUtil.checkProducerLeak; import static org.apache.flink.connector.kafka.testutils.KafkaUtil.createKafkaContainer; +import static org.apache.flink.connector.kafka.testutils.KafkaUtil.createNewTopicAndWaitForPartitionAssignment; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -176,9 +175,14 @@ class KafkaSinkITCase { } @BeforeEach - void setUp() throws ExecutionException, InterruptedException { + void setUp() { topic = UUID.randomUUID().toString(); - createTestTopic(topic, 1, TOPIC_REPLICATION_FACTOR); + Properties adminProperties = new Properties(); + adminProperties.put( + CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, + KAFKA_CONTAINER.getBootstrapServers()); + createNewTopicAndWaitForPartitionAssignment( + topic, 1, TOPIC_REPLICATION_FACTOR, adminProperties); } @AfterEach @@ -763,15 +767,6 @@ class KafkaSinkITCase { return standardProps; } - private void createTestTopic(String topic, int numPartitions, short replicationFactor) - throws ExecutionException, InterruptedException { - final CreateTopicsResult result = - admin.createTopics( - Collections.singletonList( - new NewTopic(topic, numPartitions, replicationFactor))); - result.all().get(); - } - private void deleteTestTopic(String topic) throws ExecutionException, InterruptedException { final DeleteTopicsResult result = admin.deleteTopics(Collections.singletonList(topic)); result.all().get();
