From d42dfa4c146e3037b0433312b41672d332d5d18e Mon Sep 17 00:00:00 2001 From: David Wang Date: Tue, 25 Aug 2026 13:33:50 +1000 Subject: [PATCH] Retry Kafka topic cleanup timeouts in tests --- .../cdc/kafka/KafkaActionITCaseBase.java | 48 +++++++++- .../cdc/kafka/KafkaActionITCaseBaseTest.java | 88 +++++++++++++++++++ 2 files changed, 134 insertions(+), 2 deletions(-) create mode 100644 paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBaseTest.java diff --git a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java index b2e445984e9b..553fffa11403 100644 --- a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java +++ b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBase.java @@ -26,6 +26,7 @@ import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.clients.admin.AdminClient; +import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -33,6 +34,7 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.junit.jupiter.api.AfterAll; @@ -59,6 +61,7 @@ import java.util.HashMap; import java.util.Map; import java.util.Properties; +import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; @@ -73,6 +76,9 @@ public abstract class KafkaActionITCaseBase extends CdcActionITCaseBase { private static final String INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS = "schemaregistry"; private static final Network NETWORK = Network.newNetwork(); private static final int ZK_TIMEOUT_MILLIS = 30000; + private static final int ADMIN_REQUEST_TIMEOUT_MILLIS = 60000; + private static final int DELETE_TOPICS_MAX_ATTEMPTS = 3; + private static final long DELETE_TOPICS_RETRY_BACKOFF_MILLIS = 1000L; protected static KafkaProducer kafkaProducer; private static KafkaConsumer kafkaConsumer; @@ -140,7 +146,10 @@ public static void beforeAll() { kafkaConsumer = new KafkaConsumer<>(consumerProperties); // create AdminClient - adminClient = AdminClient.create(getStandardProps()); + Properties adminProperties = getStandardProps(); + adminProperties.put( + AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, ADMIN_REQUEST_TIMEOUT_MILLIS); + adminClient = AdminClient.create(adminProperties); } @AfterAll @@ -159,7 +168,42 @@ public void after() throws Exception { } private void deleteTopics() throws ExecutionException, InterruptedException { - adminClient.deleteTopics(adminClient.listTopics().names().get()).all().get(); + retryOnTimeout( + () -> { + Set topics = adminClient.listTopics().names().get(); + if (!topics.isEmpty()) { + adminClient.deleteTopics(topics).all().get(); + } + }, + DELETE_TOPICS_MAX_ATTEMPTS, + DELETE_TOPICS_RETRY_BACKOFF_MILLIS); + } + + static void retryOnTimeout( + KafkaCleanupOperation operation, int maxAttempts, long retryBackoffMillis) + throws ExecutionException, InterruptedException { + for (int attempt = 1; attempt <= maxAttempts; attempt++) { + try { + operation.run(); + return; + } catch (ExecutionException e) { + if (!(e.getCause() instanceof TimeoutException) || attempt == maxAttempts) { + throw e; + } + + LOG.warn( + "Timed out cleaning Kafka topics on attempt {}/{}. Retrying.", + attempt, + maxAttempts, + e); + Thread.sleep(retryBackoffMillis * attempt); + } + } + } + + @FunctionalInterface + interface KafkaCleanupOperation { + void run() throws ExecutionException, InterruptedException; } public static Properties getStandardProps() { diff --git a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBaseTest.java b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBaseTest.java new file mode 100644 index 000000000000..ba3d36094952 --- /dev/null +++ b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/action/cdc/kafka/KafkaActionITCaseBaseTest.java @@ -0,0 +1,88 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.flink.action.cdc.kafka; + +import org.apache.kafka.common.errors.TimeoutException; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.ExecutionException; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests for Kafka integration test cleanup utilities. */ +class KafkaActionITCaseBaseTest { + + @Test + void testRetryKafkaTimeout() throws Exception { + AtomicInteger attempts = new AtomicInteger(); + + KafkaActionITCaseBase.retryOnTimeout( + () -> { + if (attempts.incrementAndGet() == 1) { + throw new ExecutionException(new TimeoutException("timeout")); + } + }, + 3, + 0L); + + assertThat(attempts).hasValue(2); + } + + @Test + void testDoNotRetryNonTimeoutFailure() { + AtomicInteger attempts = new AtomicInteger(); + + assertThatThrownBy( + () -> + KafkaActionITCaseBase.retryOnTimeout( + () -> { + attempts.incrementAndGet(); + throw new ExecutionException( + new IllegalStateException("failure")); + }, + 3, + 0L)) + .isInstanceOf(ExecutionException.class) + .hasCauseInstanceOf(IllegalStateException.class); + + assertThat(attempts).hasValue(1); + } + + @Test + void testFailAfterMaximumTimeoutAttempts() { + AtomicInteger attempts = new AtomicInteger(); + + assertThatThrownBy( + () -> + KafkaActionITCaseBase.retryOnTimeout( + () -> { + attempts.incrementAndGet(); + throw new ExecutionException( + new TimeoutException("timeout")); + }, + 3, + 0L)) + .isInstanceOf(ExecutionException.class) + .hasCauseInstanceOf(TimeoutException.class); + + assertThat(attempts).hasValue(3); + } +}