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();

Reply via email to