diff --git a/pulsar-client-api/pom.xml b/pulsar-client-api/pom.xml
index 6b2691866fac8..3789769cfde48 100644
--- a/pulsar-client-api/pom.xml
+++ b/pulsar-client-api/pom.xml
@@ -44,6 +44,10 @@
pulsar-transaction-common
${project.version}
+
+ io.netty
+ netty-resolver-dns
+
diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
index addedaa6c20c0..482fc9a8d935a 100644
--- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
+++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ClientBuilder.java
@@ -18,6 +18,7 @@
*/
package org.apache.pulsar.client.api;
+import io.netty.resolver.dns.DnsNameResolverBuilder;
import java.time.Clock;
import java.util.Map;
import java.util.concurrent.TimeUnit;
@@ -387,4 +388,12 @@ ClientBuilder authentication(String authPluginClassName, Map aut
* @return the client builder instance
*/
ClientBuilder clock(Clock clock);
+
+ /**
+ * The DNS resolver builder used by the pulsar client.
+ *
+ * @param builder the DNS resolver builder used by the pulsar client to resolve inet names
+ * @return the client builder instance
+ */
+ ClientBuilder dns(DnsNameResolverBuilder builder);
}
diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
index 16ea6889a9b04..470d59773a2f3 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientBuilderImpl.java
@@ -18,10 +18,10 @@
*/
package org.apache.pulsar.client.impl;
+import io.netty.resolver.dns.DnsNameResolverBuilder;
import java.time.Clock;
import java.util.Map;
import java.util.concurrent.TimeUnit;
-
import org.apache.commons.lang3.StringUtils;
import org.apache.pulsar.client.api.Authentication;
import org.apache.pulsar.client.api.AuthenticationFactory;
@@ -214,13 +214,13 @@ public ClientBuilder startingBackoffInterval(long duration, TimeUnit unit) {
conf.setInitialBackoffIntervalNanos(unit.toNanos(duration));
return this;
}
-
+
@Override
public ClientBuilder maxBackoffInterval(long duration, TimeUnit unit) {
conf.setMaxBackoffIntervalNanos(unit.toNanos(duration));
return this;
}
-
+
public ClientConfigurationData getClientConfigurationData() {
return conf;
}
@@ -230,4 +230,10 @@ public ClientBuilder clock(Clock clock) {
conf.setClock(clock);
return this;
}
+
+ @Override
+ public ClientBuilder dns(DnsNameResolverBuilder builder) {
+ conf.setDnsNameResolverBuilder(builder);
+ return this;
+ }
}
diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionPool.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionPool.java
index 30301bbd8f398..207354c1a044b 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionPool.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionPool.java
@@ -19,7 +19,6 @@
package org.apache.pulsar.client.impl;
import com.google.common.annotations.VisibleForTesting;
-
import io.netty.bootstrap.Bootstrap;
import io.netty.channel.Channel;
import io.netty.channel.ChannelException;
@@ -29,7 +28,6 @@
import io.netty.resolver.dns.DnsNameResolver;
import io.netty.resolver.dns.DnsNameResolverBuilder;
import io.netty.util.concurrent.Future;
-
import java.io.Closeable;
import java.io.IOException;
import java.net.InetAddress;
@@ -42,7 +40,6 @@
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
-
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.impl.conf.ClientConfigurationData;
import org.apache.pulsar.common.allocator.PulsarByteBufAllocator;
@@ -84,7 +81,11 @@ public ConnectionPool(ClientConfigurationData conf, EventLoopGroup eventLoopGrou
throw new PulsarClientException(e);
}
- this.dnsResolver = new DnsNameResolverBuilder(eventLoopGroup.next()).traceEnabled(true)
+ DnsNameResolverBuilder builder = conf.getDnsNameResolverBuilder();
+ if (builder == null) {
+ builder = new DnsNameResolverBuilder();
+ }
+ this.dnsResolver = builder.eventLoop(eventLoopGroup.next()).traceEnabled(true)
.channelType(EventLoopUtil.getDatagramChannelClass(eventLoopGroup)).build();
}
diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
index af478ceaf6260..4597df86c7bec 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ClientConfigurationData.java
@@ -19,7 +19,10 @@
package org.apache.pulsar.client.impl.conf;
import com.fasterxml.jackson.annotation.JsonIgnore;
+import io.netty.resolver.dns.DnsNameResolverBuilder;
+import java.io.Serializable;
import java.time.Clock;
+import java.util.concurrent.TimeUnit;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@@ -27,9 +30,6 @@
import org.apache.pulsar.client.api.ServiceUrlProvider;
import org.apache.pulsar.client.impl.auth.AuthenticationDisabled;
-import java.io.Serializable;
-import java.util.concurrent.TimeUnit;
-
/**
* This is a simple holder of the client configuration values.
*/
@@ -70,6 +70,8 @@ public class ClientConfigurationData implements Serializable, Cloneable {
private long initialBackoffIntervalNanos = TimeUnit.MILLISECONDS.toNanos(100);
private long maxBackoffIntervalNanos = TimeUnit.SECONDS.toNanos(60);
+ private DnsNameResolverBuilder dnsNameResolverBuilder = null;
+
@JsonIgnore
private Clock clock = Clock.systemDefaultZone();