From c0db83f618c3632e91705d0135afe64dcc395943 Mon Sep 17 00:00:00 2001 From: Sam Barker Date: Thu, 27 Apr 2023 12:02:07 +1200 Subject: [PATCH 1/4] Retry INVM server startup. Startup is cheap, and is usually related to transient port binding issues. Signed-off-by: Sam Barker --- .../testing/kafka/invm/InVMKafkaCluster.java | 18 ++++++++++++------ 1 file changed, 12 insertions(+), 6 deletions(-) diff --git a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java index 08d9c43e..400e3099 100644 --- a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java +++ b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java @@ -13,6 +13,7 @@ import java.net.ServerSocket; import java.nio.file.Files; import java.nio.file.Path; +import java.time.Duration; import java.util.Comparator; import java.util.List; import java.util.Map; @@ -24,13 +25,9 @@ import org.apache.kafka.common.utils.Time; import org.apache.zookeeper.server.ServerCnxnFactory; import org.apache.zookeeper.server.ZooKeeperServer; +import org.awaitility.Awaitility; import org.jetbrains.annotations.NotNull; -import io.kroxylicious.testing.kafka.api.KafkaCluster; -import io.kroxylicious.testing.kafka.common.KafkaClusterConfig; -import io.kroxylicious.testing.kafka.common.ListeningSocketPreallocator; -import io.kroxylicious.testing.kafka.common.Utils; - import kafka.server.KafkaConfig; import kafka.server.KafkaRaftServer; import kafka.server.KafkaServer; @@ -38,6 +35,11 @@ import kafka.tools.StorageTool; import scala.Option; +import io.kroxylicious.testing.kafka.api.KafkaCluster; +import io.kroxylicious.testing.kafka.common.KafkaClusterConfig; +import io.kroxylicious.testing.kafka.common.ListeningSocketPreallocator; +import io.kroxylicious.testing.kafka.common.Utils; + import static org.apache.kafka.server.common.MetadataVersion.MINIMUM_BOOTSTRAP_VERSION; /** @@ -182,7 +184,11 @@ public void start() { } } - servers.stream().parallel().forEach(Server::startup); + servers.stream().parallel().forEach(server -> Awaitility.await().atMost(Duration.ofSeconds(5)).pollDelay(Duration.ofMillis(50)).until(() -> { + //Hopefully we can remove this once a fix for https://issues.apache.org/jira/browse/KAFKA-14908 actually lands. + server.startup(); + return true; + })); Utils.awaitExpectedBrokerCountInCluster(clusterConfig.getAnonConnectConfigForCluster(kafkaEndpoints), 120, TimeUnit.SECONDS, clusterConfig.getBrokersNum()); } From a5b562a60686047c20f0b20d5a973986617e9eb7 Mon Sep 17 00:00:00 2001 From: Sam Barker Date: Thu, 27 Apr 2023 12:50:09 +1200 Subject: [PATCH 2/4] Fix formatting Signed-off-by: Sam Barker --- .../io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java index 400e3099..bd84714e 100644 --- a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java +++ b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java @@ -185,7 +185,7 @@ public void start() { } servers.stream().parallel().forEach(server -> Awaitility.await().atMost(Duration.ofSeconds(5)).pollDelay(Duration.ofMillis(50)).until(() -> { - //Hopefully we can remove this once a fix for https://issues.apache.org/jira/browse/KAFKA-14908 actually lands. + // Hopefully we can remove this once a fix for https://issues.apache.org/jira/browse/KAFKA-14908 actually lands. server.startup(); return true; })); From 4f90728e9969fe93dfb11e246129aae7a4132b37 Mon Sep 17 00:00:00 2001 From: Sam Barker Date: Thu, 27 Apr 2023 13:17:09 +1200 Subject: [PATCH 3/4] pollDelay -> pollInterval Signed-off-by: Sam Barker --- .../testing/kafka/invm/InVMKafkaCluster.java | 18 ++++++++++++------ 1 file changed, 12 insertions(+), 6 deletions(-) diff --git a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java index bd84714e..ecd5c3b4 100644 --- a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java +++ b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java @@ -22,10 +22,10 @@ import java.util.function.Supplier; import java.util.stream.Collectors; +import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.utils.Time; import org.apache.zookeeper.server.ServerCnxnFactory; import org.apache.zookeeper.server.ZooKeeperServer; -import org.awaitility.Awaitility; import org.jetbrains.annotations.NotNull; import kafka.server.KafkaConfig; @@ -41,12 +41,14 @@ import io.kroxylicious.testing.kafka.common.Utils; import static org.apache.kafka.server.common.MetadataVersion.MINIMUM_BOOTSTRAP_VERSION; +import static org.awaitility.Awaitility.await; /** * Configures and manages an in process (within the JVM) Kafka cluster. */ public class InVMKafkaCluster implements KafkaCluster { private static final System.Logger LOGGER = System.getLogger(InVMKafkaCluster.class.getName()); + private static final int STARTUP_TIMEOUT = 30; private final KafkaClusterConfig clusterConfig; private final Path tempDirectory; @@ -184,11 +186,15 @@ public void start() { } } - servers.stream().parallel().forEach(server -> Awaitility.await().atMost(Duration.ofSeconds(5)).pollDelay(Duration.ofMillis(50)).until(() -> { - // Hopefully we can remove this once a fix for https://issues.apache.org/jira/browse/KAFKA-14908 actually lands. - server.startup(); - return true; - })); + servers.stream().parallel().forEach(server -> await().atMost(Duration.ofSeconds(STARTUP_TIMEOUT)) + .catchUncaughtExceptions() + .ignoreException(KafkaException.class) + .pollInterval(Duration.ofMillis(50)) + .until(() -> { + // Hopefully we can remove this once a fix for https://issues.apache.org/jira/browse/KAFKA-14908 actually lands. + server.startup(); + return true; + })); Utils.awaitExpectedBrokerCountInCluster(clusterConfig.getAnonConnectConfigForCluster(kafkaEndpoints), 120, TimeUnit.SECONDS, clusterConfig.getBrokersNum()); } From e0abbd2e6dd834fee9681aad8b24126905200ba7 Mon Sep 17 00:00:00 2001 From: Sam Barker Date: Thu, 27 Apr 2023 21:19:21 +1200 Subject: [PATCH 4/4] Maven formatter wins again --- .../testing/kafka/invm/InVMKafkaCluster.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java index ecd5c3b4..b1166e13 100644 --- a/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java +++ b/impl/src/main/java/io/kroxylicious/testing/kafka/invm/InVMKafkaCluster.java @@ -28,6 +28,11 @@ import org.apache.zookeeper.server.ZooKeeperServer; import org.jetbrains.annotations.NotNull; +import io.kroxylicious.testing.kafka.api.KafkaCluster; +import io.kroxylicious.testing.kafka.common.KafkaClusterConfig; +import io.kroxylicious.testing.kafka.common.ListeningSocketPreallocator; +import io.kroxylicious.testing.kafka.common.Utils; + import kafka.server.KafkaConfig; import kafka.server.KafkaRaftServer; import kafka.server.KafkaServer; @@ -35,11 +40,6 @@ import kafka.tools.StorageTool; import scala.Option; -import io.kroxylicious.testing.kafka.api.KafkaCluster; -import io.kroxylicious.testing.kafka.common.KafkaClusterConfig; -import io.kroxylicious.testing.kafka.common.ListeningSocketPreallocator; -import io.kroxylicious.testing.kafka.common.Utils; - import static org.apache.kafka.server.common.MetadataVersion.MINIMUM_BOOTSTRAP_VERSION; import static org.awaitility.Awaitility.await;