From edc28a1a2e134606cd5be98e8fded7bb20dddac5 Mon Sep 17 00:00:00 2001 From: Zida Zhou Date: Fri, 3 Jun 2022 11:47:59 -0700 Subject: [PATCH] use_all_dns_ips as default value for client.dns.lookup --- .../java/org/apache/kafka/clients/ClientDnsLookup.java | 2 +- .../main/java/org/apache/kafka/clients/ClientUtils.java | 2 +- .../org/apache/kafka/clients/admin/AdminClientConfig.java | 2 +- .../org/apache/kafka/clients/consumer/ConsumerConfig.java | 2 +- .../org/apache/kafka/clients/producer/ProducerConfig.java | 2 +- .../apache/kafka/clients/admin/KafkaAdminClientTest.java | 4 ++-- .../kafka/clients/consumer/internals/FetcherTest.java | 4 ++-- .../kafka/clients/producer/internals/SenderTest.java | 2 +- .../org/apache/kafka/connect/runtime/WorkerConfig.java | 8 ++++---- .../main/scala/kafka/admin/BrokerApiVersionsCommand.scala | 4 ++-- .../scala/kafka/controller/ControllerChannelManager.scala | 2 +- .../transaction/TransactionMarkerChannelManager.scala | 2 +- core/src/main/scala/kafka/server/KafkaServer.scala | 2 +- .../scala/kafka/server/ReplicaFetcherBlockingSend.scala | 2 +- .../main/scala/kafka/tools/ReplicaVerificationTool.scala | 2 +- 15 files changed, 21 insertions(+), 21 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/ClientDnsLookup.java b/clients/src/main/java/org/apache/kafka/clients/ClientDnsLookup.java index 96d47c344a073..844f236b023d2 100644 --- a/clients/src/main/java/org/apache/kafka/clients/ClientDnsLookup.java +++ b/clients/src/main/java/org/apache/kafka/clients/ClientDnsLookup.java @@ -24,7 +24,7 @@ public enum ClientDnsLookup { USE_ALL_DNS_IPS("use_all_dns_ips"), RESOLVE_CANONICAL_BOOTSTRAP_SERVERS_ONLY("resolve_canonical_bootstrap_servers_only"); - private String clientDnsLookup; + private final String clientDnsLookup; ClientDnsLookup(String clientDnsLookup) { this.clientDnsLookup = clientDnsLookup; 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 599baeed508dd..b15060aa0a851 100644 --- a/clients/src/main/java/org/apache/kafka/clients/ClientUtils.java +++ b/clients/src/main/java/org/apache/kafka/clients/ClientUtils.java @@ -52,7 +52,7 @@ public static List parseAndValidateAddresses(List url * some third-party applications still rely on this API to parse and validate addresses. */ public static List parseAndValidateAddresses(List urls) { - return parseAndValidateAddresses(urls, ClientDnsLookup.DEFAULT); + return parseAndValidateAddresses(urls, ClientDnsLookup.USE_ALL_DNS_IPS); } public static List parseAndValidateAddresses(List urls, ClientDnsLookup clientDnsLookup) { 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 833f59b214d6a..8baa33b08cdd4 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 @@ -184,7 +184,7 @@ public class AdminClientConfig extends AbstractConfig { METRICS_RECORDING_LEVEL_DOC) .define(CLIENT_DNS_LOOKUP_CONFIG, Type.STRING, - ClientDnsLookup.DEFAULT.toString(), + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString(), ClientDnsLookup.RESOLVE_CANONICAL_BOOTSTRAP_SERVERS_ONLY.toString()), 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 e2600e27db236..acff33af8cdd9 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 @@ -328,7 +328,7 @@ public class ConsumerConfig extends AbstractConfig { CommonClientConfigs.BOOTSTRAP_SERVERS_DOC) .define(CLIENT_DNS_LOOKUP_CONFIG, Type.STRING, - ClientDnsLookup.DEFAULT.toString(), + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString(), ClientDnsLookup.RESOLVE_CANONICAL_BOOTSTRAP_SERVERS_ONLY.toString()), 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 2246ff55c4d7d..1feb7caed6c0a 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 @@ -281,7 +281,7 @@ public class ProducerConfig extends AbstractConfig { CONFIG = new ConfigDef().define(BOOTSTRAP_SERVERS_CONFIG, Type.LIST, Collections.emptyList(), new ConfigDef.NonNullValidator(), Importance.HIGH, CommonClientConfigs.BOOTSTRAP_SERVERS_DOC) .define(CLIENT_DNS_LOOKUP_CONFIG, Type.STRING, - ClientDnsLookup.DEFAULT.toString(), + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString(), ClientDnsLookup.RESOLVE_CANONICAL_BOOTSTRAP_SERVERS_ONLY.toString()), diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index 23280df57f6dc..0a2a366f6b854 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -248,7 +248,7 @@ private static Cluster mockCluster(int controllerIndex) { private static Cluster mockBootstrapCluster() { return Cluster.bootstrap(ClientUtils.parseAndValidateAddresses( - singletonList("localhost:8121"), ClientDnsLookup.DEFAULT)); + singletonList("localhost:8121"), ClientDnsLookup.USE_ALL_DNS_IPS)); } private static AdminClientUnitTestEnv mockClientEnv(String... configVals) { @@ -2281,7 +2281,7 @@ public void testGetSubLevelError() { assertEquals(FencedInstanceIdException.class, KafkaAdminClient.getSubLevelError( errorsMap, memberIdentities.get(1), "For unit test").getClass()); } - + private ClientQuotaEntity newClientQuotaEntity(String... args) { assertTrue(args.length % 2 == 0); 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 23d03b0b046ca..8ff235e07c753 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 @@ -2121,7 +2121,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, ClientDnsLookup.DEFAULT, + 1000, 1000, 64 * 1024, 64 * 1024, 1000, ClientDnsLookup.USE_ALL_DNS_IPS, time, true, new ApiVersions(), throttleTimeSensor, new LogContext()); ByteBuffer buffer = ApiVersionsResponse. @@ -3539,7 +3539,7 @@ private void testGetOffsetsForTimesWithError(Errors errorForP0, TopicPartition t2p0 = new TopicPartition(topicName2, 0); // Expect a metadata refresh. metadata.bootstrap(ClientUtils.parseAndValidateAddresses(Collections.singletonList("1.1.1.1:1111"), - ClientDnsLookup.DEFAULT)); + ClientDnsLookup.USE_ALL_DNS_IPS)); Map partitionNumByTopic = new HashMap<>(); partitionNumByTopic.put(topicName, 2); 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 1b35e0b43116e..64c372439bd51 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 @@ -271,7 +271,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, ClientDnsLookup.DEFAULT, + 1000, 1000, 64 * 1024, 64 * 1024, 1000, ClientDnsLookup.USE_ALL_DNS_IPS, time, true, new ApiVersions(), throttleTimeSensor, logContext); ByteBuffer buffer = ApiVersionsResponse.createApiVersionsResponse(400, RecordBatch.CURRENT_MAGIC_VALUE). 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 837cfe5c19c1f..263354493d1d2 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 @@ -206,7 +206,7 @@ public class WorkerConfig extends AbstractConfig { + "plugins and their dependencies\n" + "Note: symlinks will be followed to discover dependencies or plugins.\n" + "Examples: plugin.path=/usr/local/share/java,/usr/local/share/kafka/plugins," - + "/opt/connectors\n" + + "/opt/connectors\n" + "Do not use config provider variables in this property, since the raw path is used " + "by the worker's scanner before config providers are initialized and used to " + "replace variables."; @@ -249,7 +249,7 @@ protected static ConfigDef baseConfigDef() { Importance.HIGH, BOOTSTRAP_SERVERS_DOC) .define(CLIENT_DNS_LOOKUP_CONFIG, Type.STRING, - ClientDnsLookup.DEFAULT.toString(), + ClientDnsLookup.USE_ALL_DNS_IPS.toString(), in(ClientDnsLookup.DEFAULT.toString(), ClientDnsLookup.USE_ALL_DNS_IPS.toString(), ClientDnsLookup.RESOLVE_CANONICAL_BOOTSTRAP_SERVERS_ONLY.toString()), @@ -379,8 +379,8 @@ private void logPluginPathConfigProviderWarning(Map rawOriginals if (!Objects.equals(rawPluginPath, transformedPluginPath)) { log.warn( "Variables cannot be used in the 'plugin.path' property, since the property is " - + "used by plugin scanning before the config providers that replace the " - + "variables are initialized. The raw value '{}' was used for plugin scanning, as " + + "used by plugin scanning before the config providers that replace the " + + "variables are initialized. The raw value '{}' was used for plugin scanning, as " + "opposed to the transformed value '{}', and this may cause unexpected results.", rawPluginPath, transformedPluginPath diff --git a/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala b/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala index 92cdb9ee02396..2fdc3d2e687b1 100644 --- a/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala +++ b/core/src/main/scala/kafka/admin/BrokerApiVersionsCommand.scala @@ -224,7 +224,7 @@ object BrokerApiVersionsCommand { CommonClientConfigs.BOOTSTRAP_SERVERS_DOC) .define(CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG, Type.STRING, - ClientDnsLookup.DEFAULT.toString, + ClientDnsLookup.USE_ALL_DNS_IPS.toString, in(ClientDnsLookup.DEFAULT.toString, ClientDnsLookup.USE_ALL_DNS_IPS.toString, ClientDnsLookup.RESOLVE_CANONICAL_BOOTSTRAP_SERVERS_ONLY.toString), @@ -298,7 +298,7 @@ object BrokerApiVersionsCommand { DefaultSendBufferBytes, DefaultReceiveBufferBytes, requestTimeoutMs, - ClientDnsLookup.DEFAULT, + ClientDnsLookup.USE_ALL_DNS_IPS, 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 31174ef259330..8d9d446c4b270 100755 --- a/core/src/main/scala/kafka/controller/ControllerChannelManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerChannelManager.scala @@ -169,7 +169,7 @@ class ControllerChannelManager(controllerContext: ControllerContext, Selectable.USE_DEFAULT_BUFFER_SIZE, Selectable.USE_DEFAULT_BUFFER_SIZE, config.requestTimeoutMs, - ClientDnsLookup.DEFAULT, + 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 34cc6391193d3..f1843f19327e2 100644 --- a/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala +++ b/core/src/main/scala/kafka/coordinator/transaction/TransactionMarkerChannelManager.scala @@ -79,7 +79,7 @@ object TransactionMarkerChannelManager { Selectable.USE_DEFAULT_BUFFER_SIZE, config.socketReceiveBufferBytes, config.requestTimeoutMs, - ClientDnsLookup.DEFAULT, + 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 aae13d43c7232..e45257d48e340 100755 --- a/core/src/main/scala/kafka/server/KafkaServer.scala +++ b/core/src/main/scala/kafka/server/KafkaServer.scala @@ -534,7 +534,7 @@ class KafkaServer(val config: KafkaConfig, time: Time = Time.SYSTEM, threadNameP Selectable.USE_DEFAULT_BUFFER_SIZE, Selectable.USE_DEFAULT_BUFFER_SIZE, config.requestTimeoutMs, - ClientDnsLookup.DEFAULT, + 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 4f4fd1a72f150..8d647dc20e0ea 100644 --- a/core/src/main/scala/kafka/server/ReplicaFetcherBlockingSend.scala +++ b/core/src/main/scala/kafka/server/ReplicaFetcherBlockingSend.scala @@ -88,7 +88,7 @@ class ReplicaFetcherBlockingSend(sourceBroker: BrokerEndPoint, Selectable.USE_DEFAULT_BUFFER_SIZE, brokerConfig.replicaSocketReceiveBufferBytes, brokerConfig.replicaRequestTimeoutMs, - ClientDnsLookup.DEFAULT, + 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 8455b0d97e45f..c208dba55e2d7 100644 --- a/core/src/main/scala/kafka/tools/ReplicaVerificationTool.scala +++ b/core/src/main/scala/kafka/tools/ReplicaVerificationTool.scala @@ -480,7 +480,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.DEFAULT, + ClientDnsLookup.USE_ALL_DNS_IPS, time, false, new ApiVersions,