From 06d9c2a3e2e481281f293144f4f6476558450133 Mon Sep 17 00:00:00 2001 From: purshotam shah Date: Thu, 6 Aug 2026 14:04:40 -0700 Subject: [PATCH] [hotfix][runtime] Raise Kafka admin future timeout from 100ms to 30s A cold AdminClient cannot complete its first listTopics metadata round-trip in 100ms even against a local broker, so KafkaActionStateStore initialization failed nondeterministically. --- .../agents/runtime/actionstate/KafkaActionStateStore.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java index b17db1871..7c83a30e1 100644 --- a/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java +++ b/runtime/src/main/java/org/apache/flink/agents/runtime/actionstate/KafkaActionStateStore.java @@ -74,7 +74,9 @@ public class KafkaActionStateStore implements ActionStateStore { private static final Duration CONSUMER_POLL_TIMEOUT = Duration.ofMillis(1000); private static final Logger LOG = LoggerFactory.getLogger(KafkaActionStateStore.class); - private static final Long DEFAULT_FUTURE_GET_TIMEOUT_MS = 100L; + // A cold AdminClient's first metadata round-trip routinely exceeds 100ms even against a + // local broker; 100ms made store initialization fail nondeterministically at startup. + private static final Long DEFAULT_FUTURE_GET_TIMEOUT_MS = 30_000L; private final AgentConfiguration agentConfiguration;