lianetm commented on code in PR #21080:
URL: https://github.com/apache/kafka/pull/21080#discussion_r3500538515
##########
clients/src/main/java/org/apache/kafka/clients/NetworkClient.java:
##########
@@ -1210,6 +1280,139 @@ private boolean isTelemetryApi(ApiKeys apiKey) {
return apiKey == ApiKeys.GET_TELEMETRY_SUBSCRIPTIONS || apiKey ==
ApiKeys.PUSH_TELEMETRY;
}
+ public static class BootstrapConfiguration {
+ public static final BootstrapConfiguration DISABLED =
+ new BootstrapConfiguration(List.of(), null, 0, 0);
+
+ public final List<String> bootstrapServers;
+ public final ClientDnsLookup clientDnsLookup;
+ public final long bootstrapResolveTimeoutMs;
+ public final long retryBackoffMs;
+
+ private BootstrapConfiguration(final List<String> bootstrapServers,
+ final ClientDnsLookup clientDnsLookup,
+ final long bootstrapResolveTimeoutMs,
+ final long retryBackoffMs) {
+ this.bootstrapServers = bootstrapServers;
+ this.clientDnsLookup = clientDnsLookup;
+ this.bootstrapResolveTimeoutMs = bootstrapResolveTimeoutMs;
+ this.retryBackoffMs = retryBackoffMs;
+ }
+
+ public static BootstrapConfiguration enabled(final List<String>
bootstrapServers,
+ final ClientDnsLookup
clientDnsLookup,
+ final long
bootstrapResolveTimeoutMs,
+ final long
retryBackoffMs) {
+ for (String url : bootstrapServers) {
+ if (Utils.getHost(url) == null || Utils.getPort(url) == null)
+ throw new ConfigException("Invalid url in " +
CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG + ": " + url);
+ }
+ return new BootstrapConfiguration(bootstrapServers,
clientDnsLookup, bootstrapResolveTimeoutMs, retryBackoffMs);
+ }
+ }
+
+ /**
+ * Attempts to resolve bootstrap server addresses via DNS and create an
initial bootstrap cluster.
+ * This method is called from {@link #poll(long, long)} and uses a truly
asynchronous approach
+ * to avoid blocking on DNS resolution.
+ *
+ * <p>DNS resolution is performed on a separate thread via {@link
CompletableFuture}. This ensures
+ * the event loop remains responsive even if DNS lookups block or take a
long time. The bootstrap
+ * timeout can interrupt a pending DNS resolution, unlike synchronous
approaches.
+ *
+ * @param currentTimeMs The current time in milliseconds
+ * @throws BootstrapResolutionException if the bootstrap timeout expires
before DNS resolution succeeds
+ */
+ void ensureBootstrapped(final long currentTimeMs) {
+ if (bootstrapConfiguration == BootstrapConfiguration.DISABLED ||
metadataUpdater.isBootstrapped())
+ return;
+
+ if (bootstrapException != null)
+ return;
+
+ if (Thread.interrupted()) {
+ cancelBootstrapResolution();
+ throw new InterruptException(new InterruptedException());
+ }
+
+ // Check if a pending resolution completed before checking the
timeout, so that a
+ // result arriving at the same time as the deadline is not incorrectly
rejected.
+ if (pendingBootstrapResolution != null &&
pendingBootstrapResolution.isDone()) {
+ processBootstrapResolutionResult(currentTimeMs);
+ if (metadataUpdater.isBootstrapped())
+ return;
+ }
+
+ if (bootstrapTimer != null) {
+ bootstrapTimer.update(currentTimeMs);
+ checkBootstrapTimeout();
+ }
+ maybeStartBootstrapResolution(currentTimeMs);
+ }
+
+ /**
+ * Record a permanent bootstrap failure on the metadata if the timeout has
expired.
+ * The error is not thrown here; callers observe it through their metadata
layer
+ * (e.g. {@link Metadata#maybeThrowFatalException()} for
Producer/Consumer, or
+ * {@code AdminMetadataManager#bootstrapFatalException()} for AdminClient).
+ */
+ private void checkBootstrapTimeout() {
+ if (bootstrapTimer.isExpired() && bootstrapException == null) {
+ cancelBootstrapResolution();
+ bootstrapException = new BootstrapResolutionException("Failed to
resolve bootstrap servers after " +
+ bootstrapConfiguration.bootstrapResolveTimeoutMs + "ms. " +
+ "Please check your bootstrap.servers configuration and DNS
settings.");
+ metadataUpdater.bootstrapFailed(bootstrapException);
Review Comment:
since the throw was removed from here, this will now continue execution and
run the `maybeStartBootstrapResolution` that was not reached before. Will that
start a new async resolution after the failure?
Taking a step back, I wonder if we should consider reordering this, at the
moment we enforce timer then maybeStart, but the opposite seems more natural ,
wdyt? it would also simplify, no need to check bootstrapTimer != null when
enforcing the timer if we did maybeStart first. And no concerns about starting
an unneeded async resolution after the exception right? wdyt?
##########
connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedConfig.java:
##########
@@ -397,6 +397,12 @@ private static ConfigDef config(Crypto crypto) {
atLeast(0L),
ConfigDef.Importance.LOW,
CommonClientConfigs.RETRY_BACKOFF_MAX_MS_DOC)
+ .define(CommonClientConfigs.BOOTSTRAP_RESOLVE_TIMEOUT_MS_CONFIG,
+ ConfigDef.Type.LONG,
+ CommonClientConfigs.DEFAULT_BOOTSTRAP_RESOLVE_TIMEOUT_MS,
+ atLeast(0L),
+ ConfigDef.Importance.MEDIUM,
Review Comment:
this is HIGH for consumer and producer, should this align? what's the
intention?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]