diff --git a/clients/src/main/java/org/apache/kafka/clients/ClientUtils.java b/clients/src/main/java/org/apache/kafka/clients/ClientUtils.java index 0b35614801d07..63b1f14149b2e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/ClientUtils.java +++ b/clients/src/main/java/org/apache/kafka/clients/ClientUtils.java @@ -17,6 +17,7 @@ package org.apache.kafka.clients; import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.config.SaslConfigs; import org.apache.kafka.common.network.ChannelBuilder; @@ -27,8 +28,11 @@ import org.slf4j.LoggerFactory; import java.io.Closeable; +import java.net.InetAddress; import java.net.InetSocketAddress; +import java.net.UnknownHostException; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.concurrent.atomic.AtomicReference; @@ -89,4 +93,27 @@ public static ChannelBuilder createChannelBuilder(AbstractConfig config) { return ChannelBuilders.clientChannelBuilder(securityProtocol, JaasContext.Type.CLIENT, config, null, clientSaslMechanism, true); } + + static List resolve(String host, ClientDnsLookup clientDnsLookup) throws UnknownHostException { + InetAddress[] addresses = InetAddress.getAllByName(host); + if (ClientDnsLookup.USE_ALL_DNS_IPS == clientDnsLookup) { + return filterPreferredAddresses(addresses); + } else { + return Collections.singletonList(addresses[0]); + } + } + + static List filterPreferredAddresses(InetAddress[] allAddresses) { + List preferredAddresses = new ArrayList<>(); + Class clazz = null; + for (InetAddress address : allAddresses) { + if (clazz == null) { + clazz = address.getClass(); + } + if (clazz.isInstance(address)) { + preferredAddresses.add(address); + } + } + return preferredAddresses; + } } diff --git a/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java b/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java index c3a2856f4ad6b..72bad10d78cde 100644 --- a/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java +++ b/clients/src/main/java/org/apache/kafka/clients/ClusterConnectionStates.java @@ -18,9 +18,16 @@ import java.util.concurrent.ThreadLocalRandom; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.errors.AuthenticationException; +import org.apache.kafka.common.utils.LogContext; +import org.slf4j.Logger; +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; /** @@ -33,8 +40,10 @@ final class ClusterConnectionStates { private final static int RECONNECT_BACKOFF_EXP_BASE = 2; private final double reconnectBackoffMaxExp; private final Map nodeState; + private final Logger log; - public ClusterConnectionStates(long reconnectBackoffMs, long reconnectBackoffMaxMs) { + public ClusterConnectionStates(long reconnectBackoffMs, long reconnectBackoffMaxMs, LogContext logContext) { + this.log = logContext.logger(ClusterConnectionStates.class); this.reconnectBackoffInitMs = reconnectBackoffMs; this.reconnectBackoffMaxMs = reconnectBackoffMaxMs; this.reconnectBackoffMaxExp = Math.log(this.reconnectBackoffMaxMs / (double) Math.max(reconnectBackoffMs, 1)) / Math.log(RECONNECT_BACKOFF_EXP_BASE); @@ -101,25 +110,43 @@ public boolean isConnecting(String id) { } /** - * Enter the connecting state for the given connection. + * Enter the connecting state for the given connection, moving to a new resolved address if necessary. * @param id the id of the connection - * @param now the current time + * @param now the current time in ms + * @param host the host of the connection, to be resolved internally if needed + * @param clientDnsLookup the mode of DNS lookup to use when resolving the {@code host} */ - public void connecting(String id, long now) { - if (nodeState.containsKey(id)) { - NodeConnectionState node = nodeState.get(id); - node.lastConnectAttemptMs = now; - node.state = ConnectionState.CONNECTING; - } else { - nodeState.put(id, new NodeConnectionState(ConnectionState.CONNECTING, now, - this.reconnectBackoffInitMs)); + public void connecting(String id, long now, String host, ClientDnsLookup clientDnsLookup) { + NodeConnectionState connectionState = nodeState.get(id); + if (connectionState != null && connectionState.host().equals(host)) { + connectionState.lastConnectAttemptMs = now; + connectionState.state = ConnectionState.CONNECTING; + // Move to next resolved address, or if addresses are exhausted, mark node to be re-resolved + connectionState.moveToNextAddress(); + return; + } else if (connectionState != null) { + log.info("Hostname for node {} changed from {} to {}.", id, connectionState.host(), host); } + + // Create a new NodeConnectionState if nodeState does not already contain one + // for the specified id or if the hostname associated with the node id changed. + nodeState.put(id, new NodeConnectionState(ConnectionState.CONNECTING, now, + this.reconnectBackoffInitMs, host, clientDnsLookup)); + } + + /** + * Returns a resolved address for the given connection, resolving it if necessary. + * @param id the id of the connection + * @throws UnknownHostException if the address was not resolvable + */ + public InetAddress currentAddress(String id) throws UnknownHostException { + return nodeState(id).currentAddress(); } /** * Enter the disconnected state for the given node. * @param id the connection we have disconnected - * @param now the current time + * @param now the current time in ms */ public void disconnected(String id, long now) { NodeConnectionState nodeState = nodeState(id); @@ -194,7 +221,7 @@ public void ready(String id) { /** * Enter the authentication failed state for the given node. * @param id the connection identifier - * @param now the current time + * @param now the current time in ms * @param exception the authentication exception */ public void authenticationFailed(String id, long now, AuthenticationException exception) { @@ -209,7 +236,7 @@ public void authenticationFailed(String id, long now, AuthenticationException ex * Return true if the connection is in the READY state and currently not throttled. * * @param id the connection identifier - * @param now the current time + * @param now the current time in ms */ public boolean isReady(String id, long now) { return isReady(nodeState.get(id), now); @@ -223,7 +250,7 @@ private boolean isReady(NodeConnectionState state, long now) { * Return true if there is at least one node with connection in the READY state and not throttled. Returns false * otherwise. * - * @param now the current time + * @param now the current time in ms */ public boolean hasReadyNodes(long now) { for (Map.Entry entry : nodeState.entrySet()) { @@ -334,14 +361,55 @@ private static class NodeConnectionState { long reconnectBackoffMs; // Connection is being throttled if current time < throttleUntilTimeMs. long throttleUntilTimeMs; + private List addresses; + private int addressIndex; + private final String host; + private final ClientDnsLookup clientDnsLookup; - public NodeConnectionState(ConnectionState state, long lastConnectAttempt, long reconnectBackoffMs) { + private NodeConnectionState(ConnectionState state, long lastConnectAttempt, long reconnectBackoffMs, + String host, ClientDnsLookup clientDnsLookup) { this.state = state; + this.addresses = Collections.emptyList(); + this.addressIndex = -1; this.authenticationException = null; this.lastConnectAttemptMs = lastConnectAttempt; this.failedAttempts = 0; this.reconnectBackoffMs = reconnectBackoffMs; this.throttleUntilTimeMs = 0; + this.host = host; + this.clientDnsLookup = clientDnsLookup; + } + + public String host() { + return host; + } + + /** + * Fetches the current selected IP address for this node, resolving {@link #host()} if necessary. + * @return the selected address + * @throws UnknownHostException if resolving {@link #host()} fails + */ + private InetAddress currentAddress() throws UnknownHostException { + if (addresses.isEmpty()) { + // (Re-)initialize list + addresses = ClientUtils.resolve(host, clientDnsLookup); + addressIndex = 0; + } + + return addresses.get(addressIndex); + } + + /** + * Jumps to the next available resolved address for this node. If no other addresses are available, marks the + * list to be refreshed on the next {@link #currentAddress()} call. + */ + private void moveToNextAddress() { + if (addresses.isEmpty()) + return; // Avoid div0. List will initialize on next currentAddress() call + + addressIndex = (addressIndex + 1) % addresses.size(); + if (addressIndex == 0) + addresses = Collections.emptyList(); // Exhausted list. Re-resolve on next currentAddress() call } public String toString() { diff --git a/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java b/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java index f992de6b6b290..c805942230869 100644 --- a/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java +++ b/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java @@ -108,6 +108,10 @@ public class CommonClientConfigs { + "which will try to refresh metadata by choosing from existing resolved node set, this config will force resolving " + "the bootstrap url again to get new node set and use the new node set to send update metadata request"; + public static final String CLIENT_DNS_LOOKUP_CONFIG = "client.dns.lookup"; + public static final String CLIENT_DNS_LOOKUP_DOC = "

Controls how the client uses DNS lookups.

If set to use_all_dns_ips then, when the lookup returns multiple IP addresses for a hostname," + + " they will all be attempted to connect to before failing the connection. Applies to both bootstrap and advertised servers.

"; + /** * Postprocess the configuration so that exponential backoff is disabled when reconnect backoff * is explicitly configured but the maximum reconnect backoff is not explicitly configured. diff --git a/clients/src/main/java/org/apache/kafka/clients/Metadata.java b/clients/src/main/java/org/apache/kafka/clients/Metadata.java index a95b730ce73f6..ce8f67a6f5a56 100644 --- a/clients/src/main/java/org/apache/kafka/clients/Metadata.java +++ b/clients/src/main/java/org/apache/kafka/clients/Metadata.java @@ -155,8 +155,9 @@ public synchronized void incrementNodesTriedSinceLastSuccessfulRefresh() { * has been set by receiving stale metadata from a different cluster */ public synchronized boolean shouldUpdateClusterMetadataFromBootstrap(long nowMs) { - return (this.nodesTriedSinceLastSuccessfulRefresh >= 1 && - this.lastSuccessfulRefreshMs + this.maxClusterMetadataExpireTimeMs <= nowMs) || + return this.maxClusterMetadataExpireTimeMs > 0 && + (this.nodesTriedSinceLastSuccessfulRefresh >= 1 && + (this.lastRefreshMs != 0 && this.maxClusterMetadataExpireTimeMs <= nowMs - this.lastSuccessfulRefreshMs)) || this.forceClusterMetadataUpdateFromBootstrap; } diff --git a/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java b/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java index 63fcf390328ff..9f095fbee4573 100644 --- a/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/NetworkClient.java @@ -20,6 +20,7 @@ import org.apache.kafka.common.Cluster; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.errors.AuthenticationException; import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.memory.MemoryPool; @@ -46,6 +47,7 @@ import org.slf4j.Logger; import java.io.IOException; +import java.net.InetAddress; import java.net.InetSocketAddress; import java.nio.ByteBuffer; import java.util.ArrayList; @@ -100,6 +102,8 @@ public class NetworkClient implements KafkaClient { /* time in ms to wait before retrying to create connection to a server */ private final long reconnectBackoffMs; + private final ClientDnsLookup clientDnsLookup; + private final Time time; private boolean enableStickyMetadataFetch = true; @@ -130,6 +134,7 @@ public NetworkClient(Selectable selector, int socketSendBuffer, int socketReceiveBuffer, int defaultRequestTimeoutMs, + ClientDnsLookup clientDnsLookup, Time time, boolean discoverBrokerVersions, ApiVersions apiVersions, @@ -144,6 +149,7 @@ public NetworkClient(Selectable selector, socketSendBuffer, socketReceiveBuffer, defaultRequestTimeoutMs, + clientDnsLookup, time, discoverBrokerVersions, apiVersions, @@ -161,6 +167,7 @@ public NetworkClient(Selectable selector, int socketSendBuffer, int socketReceiveBuffer, int defaultRequestTimeoutMs, + ClientDnsLookup clientDnsLookup, Time time, boolean discoverBrokerVersions, ApiVersions apiVersions, @@ -177,6 +184,7 @@ public NetworkClient(Selectable selector, socketSendBuffer, socketReceiveBuffer, defaultRequestTimeoutMs, + clientDnsLookup, time, discoverBrokerVersions, apiVersions, @@ -194,6 +202,7 @@ public NetworkClient(Selectable selector, int socketSendBuffer, int socketReceiveBuffer, int defaultRequestTimeoutMs, + ClientDnsLookup clientDnsLookup, Time time, boolean discoverBrokerVersions, ApiVersions apiVersions, @@ -208,6 +217,7 @@ public NetworkClient(Selectable selector, socketSendBuffer, socketReceiveBuffer, defaultRequestTimeoutMs, + clientDnsLookup, time, discoverBrokerVersions, apiVersions, @@ -225,6 +235,7 @@ public NetworkClient(Selectable selector, int socketSendBuffer, int socketReceiveBuffer, int defaultRequestTimeoutMs, + ClientDnsLookup clientDnsLookup, Time time, boolean discoverBrokerVersions, ApiVersions apiVersions, @@ -240,6 +251,7 @@ public NetworkClient(Selectable selector, socketSendBuffer, socketReceiveBuffer, defaultRequestTimeoutMs, + clientDnsLookup, time, discoverBrokerVersions, apiVersions, @@ -258,6 +270,7 @@ private NetworkClient(MetadataUpdater metadataUpdater, int socketSendBuffer, int socketReceiveBuffer, int defaultRequestTimeoutMs, + ClientDnsLookup clientDnsLookup, Time time, boolean discoverBrokerVersions, ApiVersions apiVersions, @@ -278,7 +291,7 @@ private NetworkClient(MetadataUpdater metadataUpdater, this.selector = selector; this.clientId = clientId; this.inFlightRequests = new InFlightRequests(maxInFlightRequestsPerConnection); - this.connectionStates = new ClusterConnectionStates(reconnectBackoffMs, reconnectBackoffMax); + this.connectionStates = new ClusterConnectionStates(reconnectBackoffMs, reconnectBackoffMax, logContext); this.socketSendBuffer = socketSendBuffer; this.socketReceiveBuffer = socketReceiveBuffer; this.correlation = 0; @@ -291,6 +304,7 @@ private NetworkClient(MetadataUpdater metadataUpdater, this.throttleTimeSensor = throttleTimeSensor; this.logContext = logContext; this.log = logContext.logger(NetworkClient.class); + this.clientDnsLookup = clientDnsLookup; this.bootstrapServers.addAll(bootstrapServersConfig); } @@ -952,22 +966,25 @@ private static void correlate(RequestHeader requestHeader, ResponseHeader respon /** * Initiate a connection to the given node + * @param node the node to connect to + * @param now current time in epoch milliseconds */ private void initiateConnect(Node node, long now) { String nodeConnectionId = node.idString(); try { - log.debug("Initiating connection to node {}", node); - this.connectionStates.connecting(nodeConnectionId, now); + connectionStates.connecting(nodeConnectionId, now, node.host(), clientDnsLookup); + InetAddress address = connectionStates.currentAddress(nodeConnectionId); + log.debug("Initiating connection to node {} using address {}", node, address); selector.connect(nodeConnectionId, - new InetSocketAddress(node.host(), node.port()), - this.socketSendBuffer, - this.socketReceiveBuffer); + new InetSocketAddress(address, node.port()), + this.socketSendBuffer, + this.socketReceiveBuffer); } catch (IOException e) { + log.warn("Error connecting to node {}", node, e); /* attempt failed, we'll try again after the backoff */ connectionStates.disconnected(nodeConnectionId, now); /* maybe the problem is our metadata, update it */ metadataUpdater.requestUpdate(); - log.warn("Error connecting to node {}", node, e); } } diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java b/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java index f29413239ba13..36d0daed5f310 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java @@ -19,6 +19,7 @@ import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; @@ -165,6 +166,13 @@ public class AdminClientConfig extends AbstractConfig { false, Importance.LOW, METRICS_REPLACE_ON_DUPLICATE_DOC) + .define(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG, + Type.STRING, + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), + in(ClientDnsLookup.DEFAULT.toString(), + ClientDnsLookup.USE_ALL_DNS_IPS.toString()), + Importance.MEDIUM, + CommonClientConfigs.CLIENT_DNS_LOOKUP_DOC) // security support .define(SECURITY_PROTOCOL_CONFIG, Type.STRING, diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index ab77b504489ba..ba12f41fe0abf 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -21,6 +21,7 @@ import org.apache.kafka.clients.ClientRequest; import org.apache.kafka.clients.ClientResponse; import org.apache.kafka.clients.ClientUtils; +import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.clients.KafkaClient; import org.apache.kafka.clients.NetworkClient; import org.apache.kafka.clients.StaleMetadataException; @@ -45,6 +46,7 @@ import org.apache.kafka.common.acl.AclBinding; import org.apache.kafka.common.acl.AclBindingFilter; import org.apache.kafka.common.annotation.InterfaceStability; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.AuthenticationException; @@ -359,6 +361,7 @@ static KafkaAdminClient createInternal(AdminClientConfig config, TimeoutProcesso config.getInt(AdminClientConfig.SEND_BUFFER_CONFIG), config.getInt(AdminClientConfig.RECEIVE_BUFFER_CONFIG), (int) TimeUnit.HOURS.toMillis(1), + ClientDnsLookup.forConfig(config.getString(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG)), time, true, apiVersions, diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java index 1a424ff2da6b4..363e505c03fc0 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java @@ -18,6 +18,7 @@ import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; @@ -500,6 +501,12 @@ public class ConsumerConfig extends AbstractConfig { DEFAULT_ALLOW_AUTO_CREATE_TOPICS, Importance.MEDIUM, ALLOW_AUTO_CREATE_TOPICS_DOC) + .define(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG, + Type.STRING, + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), + in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString()), + Importance.MEDIUM, + CommonClientConfigs.CLIENT_DNS_LOOKUP_DOC) // security support .define(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, Type.STRING, diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java index 1eee9fc6af1e5..c1eea5e4182ad 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java @@ -36,6 +36,7 @@ import org.apache.kafka.common.MetricName; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.errors.InterruptException; import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.internals.ClusterResourceListeners; @@ -738,6 +739,7 @@ private KafkaConsumer(ConsumerConfig config, config.getInt(ConsumerConfig.SEND_BUFFER_CONFIG), config.getInt(ConsumerConfig.RECEIVE_BUFFER_CONFIG), config.getInt(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG), + ClientDnsLookup.forConfig(config.getString(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG)), time, true, new ApiVersions(), diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java index dc01cf76b169d..d582a09ee2940 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java @@ -49,6 +49,7 @@ import org.apache.kafka.common.MetricName; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.AuthenticationException; @@ -425,7 +426,9 @@ public KafkaProducer(Properties properties, Serializer keySerializer, Seriali config.getLong(ProducerConfig.RECONNECT_BACKOFF_MS_CONFIG), config.getLong(ProducerConfig.RECONNECT_BACKOFF_MAX_MS_CONFIG), config.getInt(ProducerConfig.SEND_BUFFER_CONFIG), - config.getInt(ProducerConfig.RECEIVE_BUFFER_CONFIG), this.requestTimeoutMs, time, true, apiVersions, + config.getInt(ProducerConfig.RECEIVE_BUFFER_CONFIG), this.requestTimeoutMs, + ClientDnsLookup.forConfig(producerConfig.getString(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG)), + time, true, apiVersions, throttleTimeSensor, logContext, producerConfig.getList(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG)); networkClient.setEnableStickyMetadataFetch(config.getBoolean(CommonClientConfigs.ENABLE_STICKY_METADATA_FETCH_CONFIG)); client = networkClient; diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java index dfdf8789f935b..19ddecd502512 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java @@ -20,6 +20,7 @@ import org.apache.kafka.clients.Metadata; import org.apache.kafka.clients.producer.internals.DefaultPartitioner; import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; @@ -380,7 +381,13 @@ public class ProducerConfig extends AbstractConfig { Type.BOOLEAN, DEFAULT_ALLOW_AUTO_CREATE_TOPICS, Importance.MEDIUM, - ALLOW_AUTO_CREATE_TOPICS_DOC); + ALLOW_AUTO_CREATE_TOPICS_DOC) + .define(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG, + Type.STRING, + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), + in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString()), + Importance.MEDIUM, + CommonClientConfigs.CLIENT_DNS_LOOKUP_DOC); } @Override diff --git a/clients/src/main/java/org/apache/kafka/common/config/ClientDnsLookup.java b/clients/src/main/java/org/apache/kafka/common/config/ClientDnsLookup.java new file mode 100644 index 0000000000000..2186b64230f52 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/common/config/ClientDnsLookup.java @@ -0,0 +1,40 @@ +/* + * 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.kafka.common.config; + +import java.util.Locale; + +public enum ClientDnsLookup { + + DEFAULT("default"), + USE_ALL_DNS_IPS("use_all_dns_ips"); + + private final String clientDnsLookup; + + ClientDnsLookup(String clientDnsLookup) { + this.clientDnsLookup = clientDnsLookup; + } + + @Override + public String toString() { + return clientDnsLookup; + } + + public static ClientDnsLookup forConfig(String config) { + return ClientDnsLookup.valueOf(config.toUpperCase(Locale.ROOT)); + } +} diff --git a/clients/src/test/java/org/apache/kafka/clients/ClientUtilsTest.java b/clients/src/test/java/org/apache/kafka/clients/ClientUtilsTest.java index d19b0be4be49f..3f21208f691bd 100644 --- a/clients/src/test/java/org/apache/kafka/clients/ClientUtilsTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/ClientUtilsTest.java @@ -16,11 +16,17 @@ */ package org.apache.kafka.clients; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigException; +import org.junit.Ignore; import org.junit.Test; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import java.net.InetAddress; import java.net.InetSocketAddress; +import java.net.UnknownHostException; import java.util.Arrays; import java.util.List; @@ -49,6 +55,41 @@ public void testOnlyBadHostname() { check("some.invalid.hostname.foo.bar.local:9999"); } + @Test + public void testFilterPreferredAddresses() throws UnknownHostException { + InetAddress ipv4 = InetAddress.getByName("192.0.0.1"); + InetAddress ipv6 = InetAddress.getByName("::1"); + + InetAddress[] ipv4First = new InetAddress[]{ipv4, ipv6, ipv4}; + List result = ClientUtils.filterPreferredAddresses(ipv4First); + assertTrue(result.contains(ipv4)); + assertFalse(result.contains(ipv6)); + assertEquals(2, result.size()); + + InetAddress[] ipv6First = new InetAddress[]{ipv6, ipv4, ipv4}; + result = ClientUtils.filterPreferredAddresses(ipv6First); + assertTrue(result.contains(ipv6)); + assertFalse(result.contains(ipv4)); + assertEquals(1, result.size()); + } + + @Test(expected = UnknownHostException.class) + public void testResolveUnknownHostException() throws UnknownHostException { + ClientUtils.resolve("some.invalid.hostname.foo.bar.local", ClientDnsLookup.USE_ALL_DNS_IPS); + } + + @Test + public void testResolveDnsLookup() throws UnknownHostException { + assertEquals(1, ClientUtils.resolve("localhost", ClientDnsLookup.DEFAULT).size()); + } + + @Test + @Ignore + public void testResolveDnsLookupAllIps() throws UnknownHostException { + // Note that kafka.apache.org resolves to 2 IP addresses + assertEquals(2, ClientUtils.resolve("kafka.apache.org", ClientDnsLookup.USE_ALL_DNS_IPS).size()); + } + private List check(String... url) { return ClientUtils.parseAndValidateAddresses(Arrays.asList(url)); } diff --git a/clients/src/test/java/org/apache/kafka/clients/ClusterConnectionStatesTest.java b/clients/src/test/java/org/apache/kafka/clients/ClusterConnectionStatesTest.java index 37155ce8ef8da..98dbd638b3235 100644 --- a/clients/src/test/java/org/apache/kafka/clients/ClusterConnectionStatesTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/ClusterConnectionStatesTest.java @@ -19,12 +19,20 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotSame; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; +import java.net.InetAddress; +import java.net.UnknownHostException; + +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.errors.AuthenticationException; +import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.MockTime; import org.junit.Before; +import org.junit.Ignore; import org.junit.Test; public class ClusterConnectionStatesTest { @@ -35,12 +43,13 @@ public class ClusterConnectionStatesTest { private final double reconnectBackoffJitter = 0.2; private final String nodeId1 = "1001"; private final String nodeId2 = "2002"; + private final String hostTwoIps = "kafka.apache.org"; private ClusterConnectionStates connectionStates; @Before public void setup() { - this.connectionStates = new ClusterConnectionStates(reconnectBackoffMs, reconnectBackoffMax); + this.connectionStates = new ClusterConnectionStates(reconnectBackoffMs, reconnectBackoffMax, new LogContext()); } @Test @@ -48,7 +57,7 @@ public void testClusterConnectionStateChanges() { assertTrue(connectionStates.canConnect(nodeId1, time.milliseconds())); // Start connecting to Node and check state - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); assertEquals(connectionStates.connectionState(nodeId1), ConnectionState.CONNECTING); assertTrue(connectionStates.isConnecting(nodeId1)); assertFalse(connectionStates.isReady(nodeId1, time.milliseconds())); @@ -96,7 +105,7 @@ public void testMultipleNodeConnectionStates() { // Start connecting one node and check that the pool only shows ready nodes after // successful connect - connectionStates.connecting(nodeId2, time.milliseconds()); + connectionStates.connecting(nodeId2, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); assertFalse(connectionStates.hasReadyNodes(time.milliseconds())); time.sleep(1000); connectionStates.ready(nodeId2); @@ -104,7 +113,7 @@ public void testMultipleNodeConnectionStates() { // Connect second node and check that both are shown as ready, pool should immediately // show ready nodes, since node2 is already connected - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); assertTrue(connectionStates.hasReadyNodes(time.milliseconds())); time.sleep(1000); connectionStates.ready(nodeId1); @@ -128,7 +137,7 @@ public void testMultipleNodeConnectionStates() { @Test public void testAuthorizationFailed() { // Try connecting - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); time.sleep(100); @@ -148,7 +157,7 @@ public void testAuthorizationFailed() { @Test public void testRemoveNode() { - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); time.sleep(1000); connectionStates.ready(nodeId1); time.sleep(10000); @@ -164,7 +173,7 @@ public void testRemoveNode() { @Test public void testMaxReconnectBackoff() { long effectiveMaxReconnectBackoff = Math.round(reconnectBackoffMax * (1 + reconnectBackoffJitter)); - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); time.sleep(1000); connectionStates.disconnected(nodeId1, time.milliseconds()); @@ -175,7 +184,7 @@ public void testMaxReconnectBackoff() { assertFalse(connectionStates.canConnect(nodeId1, time.milliseconds())); time.sleep(reconnectBackoff + 1); assertTrue(connectionStates.canConnect(nodeId1, time.milliseconds())); - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); time.sleep(10); connectionStates.disconnected(nodeId1, time.milliseconds()); } @@ -190,7 +199,7 @@ public void testExponentialReconnectBackoff() { // Run through 10 disconnects and check that reconnect backoff value is within expected range for every attempt for (int i = 0; i < 10; i++) { - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); connectionStates.disconnected(nodeId1, time.milliseconds()); // Calculate expected backoff value without jitter long expectedBackoff = Math.round(Math.pow(reconnectBackoffExpBase, Math.min(i, reconnectBackoffMaxExp)) @@ -203,7 +212,7 @@ public void testExponentialReconnectBackoff() { @Test public void testThrottled() { - connectionStates.connecting(nodeId1, time.milliseconds()); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); time.sleep(1000); connectionStates.ready(nodeId1); time.sleep(10000); @@ -226,4 +235,49 @@ public void testThrottled() { assertEquals(connectionStates.connectionDelay(nodeId1, time.milliseconds()), connectionStates.pollDelayMs(nodeId1, time.milliseconds())); } + + @Test + public void testSingleIPWithDefault() throws UnknownHostException { + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); + InetAddress currAddress = connectionStates.currentAddress(nodeId1); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.DEFAULT); + assertSame(currAddress, connectionStates.currentAddress(nodeId1)); + } + + @Test + public void testSingleIPWithUseAll() throws UnknownHostException { + assertEquals(1, ClientUtils.resolve("localhost", ClientDnsLookup.USE_ALL_DNS_IPS).size()); + + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.USE_ALL_DNS_IPS); + InetAddress currAddress = connectionStates.currentAddress(nodeId1); + connectionStates.connecting(nodeId1, time.milliseconds(), "localhost", ClientDnsLookup.USE_ALL_DNS_IPS); + assertSame(currAddress, connectionStates.currentAddress(nodeId1)); + } + + @Ignore + @Test + public void testMultipleIPsWithDefault() throws UnknownHostException { + assertEquals(2, ClientUtils.resolve(hostTwoIps, ClientDnsLookup.USE_ALL_DNS_IPS).size()); + + connectionStates.connecting(nodeId1, time.milliseconds(), hostTwoIps, ClientDnsLookup.DEFAULT); + InetAddress currAddress = connectionStates.currentAddress(nodeId1); + connectionStates.connecting(nodeId1, time.milliseconds(), hostTwoIps, ClientDnsLookup.DEFAULT); + assertSame(currAddress, connectionStates.currentAddress(nodeId1)); + } + + @Ignore + @Test + public void testMultipleIPsWithUseAll() throws UnknownHostException { + assertEquals(2, ClientUtils.resolve(hostTwoIps, ClientDnsLookup.USE_ALL_DNS_IPS).size()); + + connectionStates.connecting(nodeId1, time.milliseconds(), hostTwoIps, ClientDnsLookup.USE_ALL_DNS_IPS); + InetAddress addr1 = connectionStates.currentAddress(nodeId1); + connectionStates.connecting(nodeId1, time.milliseconds(), hostTwoIps, ClientDnsLookup.USE_ALL_DNS_IPS); + InetAddress addr2 = connectionStates.currentAddress(nodeId1); + assertNotSame(addr1, addr2); + + connectionStates.connecting(nodeId1, time.milliseconds(), hostTwoIps, ClientDnsLookup.USE_ALL_DNS_IPS); + InetAddress addr3 = connectionStates.currentAddress(nodeId1); + assertSame(addr1, addr3); + } } diff --git a/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java b/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java index a51e3d8e8d2e4..dbd5ee6a356f5 100644 --- a/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/NetworkClientTest.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.Cluster; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.network.NetworkReceive; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.CommonFields; @@ -73,25 +74,26 @@ public class NetworkClientTest { private NetworkClient createNetworkClient(long reconnectBackoffMaxMs) { return new NetworkClient(selector, metadata, "mock", Integer.MAX_VALUE, reconnectBackoffMsTest, reconnectBackoffMaxMs, 64 * 1024, 64 * 1024, - minRequestTimeoutMs, time, true, new ApiVersions(), new LogContext()); + minRequestTimeoutMs, ClientDnsLookup.DEFAULT, time, true, new ApiVersions(), new LogContext()); } private NetworkClient createNetworkClientWithStaticNodes() { return new NetworkClient(selector, new ManualMetadataUpdater(Arrays.asList(node)), "mock-static", Integer.MAX_VALUE, 0, 0, 64 * 1024, 64 * 1024, minRequestTimeoutMs, - time, true, new ApiVersions(), new LogContext()); + ClientDnsLookup.DEFAULT, time, true, new ApiVersions(), new LogContext()); } private NetworkClient createNetworkClientWithNoVersionDiscovery() { return new NetworkClient(selector, metadata, "mock", Integer.MAX_VALUE, reconnectBackoffMsTest, reconnectBackoffMaxMsTest, - 64 * 1024, 64 * 1024, minRequestTimeoutMs, time, false, new ApiVersions(), new LogContext()); + 64 * 1024, 64 * 1024, minRequestTimeoutMs, ClientDnsLookup.DEFAULT, + time, false, new ApiVersions(), new LogContext()); } private NetworkClient createClusterNetworkClient() { return new NetworkClient(selector, clusterMetadataUpdater, "mock-cluster-md", Integer.MAX_VALUE, 0, 0, 64 * 1024, 64 * 1024, - minRequestTimeoutMs, time, true, new ApiVersions(), new LogContext(), + minRequestTimeoutMs, ClientDnsLookup.DEFAULT, time, true, new ApiVersions(), new LogContext(), Collections.singletonList("example.com:10000")); } @@ -132,6 +134,12 @@ public void testSimpleRequestResponseWithNoBrokerDiscovery() { checkSimpleRequestResponse(clientWithNoVersionDiscovery); } + @Test + public void testDnsLookupFailure() { + /* Fail cleanly when the node has a bad hostname */ + assertFalse(client.ready(new Node(1234, "badhost", 1234), time.milliseconds())); + } + @Test public void testClose() { client.ready(node, time.milliseconds()); diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java index 9389ba743b0a2..01f7e59e2894c 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/FetcherTest.java @@ -39,6 +39,7 @@ import org.apache.kafka.common.Node; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.errors.InvalidTopicException; import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.errors.SerializationException; @@ -1626,7 +1627,7 @@ public void testQuotaMetrics() { Cluster cluster = TestUtils.singletonCluster("test", 1); Node node = cluster.nodes().get(0); NetworkClient client = new NetworkClient(selector, metadata, "mock", Integer.MAX_VALUE, - 1000, 1000, 64 * 1024, 64 * 1024, 1000, + 1000, 1000, 64 * 1024, 64 * 1024, 1000, ClientDnsLookup.USE_ALL_DNS_IPS, time, true, new ApiVersions(), throttleTimeSensor, new LogContext(), Collections.emptyList()); short apiVersionsResponseVersion = ApiKeys.API_VERSIONS.latestVersion(); diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java index 6dd466f8afdbe..67a461463daf9 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java @@ -44,6 +44,7 @@ import org.apache.kafka.common.MetricNameTemplate; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.errors.ClusterAuthorizationException; import org.apache.kafka.common.errors.NetworkException; import org.apache.kafka.common.errors.OutOfOrderSequenceException; @@ -262,7 +263,7 @@ public void testQuotaMetrics() throws Exception { Cluster cluster = TestUtils.singletonCluster("test", 1); Node node = cluster.nodes().get(0); NetworkClient client = new NetworkClient(selector, metadata, "mock", Integer.MAX_VALUE, - 1000, 1000, 64 * 1024, 64 * 1024, 1000, + 1000, 1000, 64 * 1024, 64 * 1024, 1000, ClientDnsLookup.USE_ALL_DNS_IPS, time, true, new ApiVersions(), throttleTimeSensor, logContext, Collections.emptyList()); short apiVersionsResponseVersion = ApiKeys.API_VERSIONS.latestVersion(); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 583953d7f591f..2ba406b83bb4e 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -18,6 +18,7 @@ import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.common.config.AbstractConfig; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; @@ -274,7 +275,13 @@ protected static ConfigDef baseConfigDef() { Collections.emptyList(), Importance.LOW, CONFIG_PROVIDERS_DOC) .define(REST_EXTENSION_CLASSES_CONFIG, Type.LIST, "", - Importance.LOW, REST_EXTENSION_CLASSES_DOC); + Importance.LOW, REST_EXTENSION_CLASSES_DOC) + .define(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG, + Type.STRING, + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), + in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString()), + Importance.MEDIUM, + CommonClientConfigs.CLIENT_DNS_LOOKUP_DOC); } private void logInternalConverterDeprecationWarnings(Map props) { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java index 525ce7e2de93d..16206fbe9781d 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java @@ -24,6 +24,7 @@ import org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.KafkaException; +import org.apache.kafka.common.config.ClientDnsLookup; import org.apache.kafka.common.metrics.JmxReporter; import org.apache.kafka.common.metrics.MetricConfig; import org.apache.kafka.common.metrics.Metrics; @@ -107,6 +108,7 @@ public WorkerGroupMember(DistributedConfig config, config.getInt(CommonClientConfigs.SEND_BUFFER_CONFIG), config.getInt(CommonClientConfigs.RECEIVE_BUFFER_CONFIG), config.getInt(CommonClientConfigs.REQUEST_TIMEOUT_MS_CONFIG), + ClientDnsLookup.forConfig(config.getString(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG)), time, true, new ApiVersions(), diff --git a/core/src/main/scala/kafka/admin/AdminClient.scala b/core/src/main/scala/kafka/admin/AdminClient.scala index 239844d2f5940..9342d4075ae2b 100644 --- a/core/src/main/scala/kafka/admin/AdminClient.scala +++ b/core/src/main/scala/kafka/admin/AdminClient.scala @@ -39,6 +39,7 @@ import org.apache.kafka.common.{Cluster, Node, TopicPartition} import scala.collection.JavaConverters._ import scala.util.{Failure, Success, Try} +import org.apache.kafka.common.config.ClientDnsLookup /** * A Scala administrative client for Kafka which supports managing and inspecting topics, brokers, @@ -452,6 +453,7 @@ object AdminClient { DefaultSendBufferBytes, DefaultReceiveBufferBytes, requestTimeoutMs, + ClientDnsLookup.USE_ALL_DNS_IPS, time, true, new ApiVersions, diff --git a/core/src/main/scala/kafka/consumer/SSLNetworkClient.scala b/core/src/main/scala/kafka/consumer/SSLNetworkClient.scala index 7d96cddc37dc3..5adcd4742f928 100644 --- a/core/src/main/scala/kafka/consumer/SSLNetworkClient.scala +++ b/core/src/main/scala/kafka/consumer/SSLNetworkClient.scala @@ -22,6 +22,7 @@ import java.util.concurrent.TimeUnit import org.apache.kafka.clients._ import org.apache.kafka.common.Node +import org.apache.kafka.common.config.ClientDnsLookup import org.apache.kafka.common.metrics.{JmxReporter, MetricConfig, Metrics, MetricsReporter} import org.apache.kafka.common.network.{ChannelBuilders, NetworkReceive, Selectable, Selector => KSelector} import org.apache.kafka.common.protocol.ApiKeys @@ -77,6 +78,7 @@ class SSLNetworkClient(config: ConsumerConfig, metadataUpdater: ManualMetadataUp Selectable.USE_DEFAULT_BUFFER_SIZE, config.socketReceiveBufferBytes, socketTimeoutMs, + ClientDnsLookup.DEFAULT, time, true, new ApiVersions, diff --git a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala index 6d8e3c3e40c25..844aac0186325 100755 --- a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala @@ -39,6 +39,7 @@ import org.apache.kafka.common.{KafkaException, Node, TopicPartition} import scala.collection.JavaConverters._ import scala.collection.mutable.HashMap import scala.collection.{Map, Set, mutable} +import org.apache.kafka.common.config.ClientDnsLookup object ControllerChannelManager { @@ -155,6 +156,7 @@ class ControllerChannelManager(controllerContext: ControllerContext, config: Kaf Selectable.USE_DEFAULT_BUFFER_SIZE, Selectable.USE_DEFAULT_BUFFER_SIZE, config.requestTimeoutMs, + ClientDnsLookup.USE_ALL_DNS_IPS, time, false, new ApiVersions, diff --git a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala index f8b56e8e00800..af9261cbdc976 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala @@ -36,6 +36,7 @@ import java.util.concurrent.{BlockingQueue, ConcurrentHashMap, LinkedBlockingQue import collection.JavaConverters._ import scala.collection.{concurrent, immutable} +import org.apache.kafka.common.config.ClientDnsLookup object TransactionMarkerChannelManager { def apply(config: KafkaConfig, @@ -74,6 +75,7 @@ object TransactionMarkerChannelManager { Selectable.USE_DEFAULT_BUFFER_SIZE, config.socketReceiveBufferBytes, config.requestTimeoutMs, + ClientDnsLookup.USE_ALL_DNS_IPS, time, false, new ApiVersions, diff --git a/core/src/main/scala/kafka/server/KafkaServer.scala b/core/src/main/scala/kafka/server/KafkaServer.scala index 338870711e2f9..633851cdb6d40 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -51,6 +51,7 @@ import org.apache.kafka.common.{ClusterResource, Node} import scala.collection.JavaConverters._ import scala.collection.{Map, Seq, mutable} +import org.apache.kafka.common.config.ClientDnsLookup object KafkaServer { // Copy the subset of properties that are relevant to Logs @@ -472,6 +473,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP Selectable.USE_DEFAULT_BUFFER_SIZE, Selectable.USE_DEFAULT_BUFFER_SIZE, config.requestTimeoutMs, + ClientDnsLookup.USE_ALL_DNS_IPS, time, false, new ApiVersions, diff --git a/core/src/main/scala/kafka/server/ReplicaFetcherBlockingSend.scala b/core/src/main/scala/kafka/server/ReplicaFetcherBlockingSend.scala index 4c7adfbe9cd5f..6027a57a2563d 100644 --- a/core/src/main/scala/kafka/server/ReplicaFetcherBlockingSend.scala +++ b/core/src/main/scala/kafka/server/ReplicaFetcherBlockingSend.scala @@ -30,6 +30,7 @@ import org.apache.kafka.common.Node import org.apache.kafka.common.requests.AbstractRequest.Builder import scala.collection.JavaConverters._ +import org.apache.kafka.common.config.ClientDnsLookup trait BlockingSend { @@ -79,6 +80,7 @@ class ReplicaFetcherBlockingSend(sourceBroker: BrokerEndPoint, Selectable.USE_DEFAULT_BUFFER_SIZE, brokerConfig.replicaSocketReceiveBufferBytes, brokerConfig.requestTimeoutMs, + ClientDnsLookup.USE_ALL_DNS_IPS, time, false, new ApiVersions, diff --git a/core/src/main/scala/kafka/tools/ReplicaVerificationTool.scala b/core/src/main/scala/kafka/tools/ReplicaVerificationTool.scala index 1ecea09b9946d..d5512871436d1 100644 --- a/core/src/main/scala/kafka/tools/ReplicaVerificationTool.scala +++ b/core/src/main/scala/kafka/tools/ReplicaVerificationTool.scala @@ -43,6 +43,7 @@ import org.apache.kafka.common.utils.{LogContext, Time} import org.apache.kafka.common.{Node, TopicPartition} import scala.collection.JavaConverters._ +import org.apache.kafka.common.config.ClientDnsLookup /** * For verifying the consistency among replicas. @@ -471,6 +472,7 @@ private class ReplicaFetcherBlockingSend(sourceNode: Node, Selectable.USE_DEFAULT_BUFFER_SIZE, consumerConfig.getInt(ConsumerConfig.RECEIVE_BUFFER_CONFIG), consumerConfig.getInt(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG), + ClientDnsLookup.USE_ALL_DNS_IPS, time, false, new ApiVersions,