From 5192b66e4e559a9ee71b08d1c75da551b600f850 Mon Sep 17 00:00:00 2001 From: Diego Salvi Date: Thu, 20 Feb 2020 15:01:03 +0100 Subject: [PATCH] Add configurable DNS client for ConnectionPool --- pulsar-client-api/pom.xml | 4 ++++ .../org/apache/pulsar/client/api/ClientBuilder.java | 9 +++++++++ .../apache/pulsar/client/impl/ClientBuilderImpl.java | 12 +++++++++--- .../apache/pulsar/client/impl/ConnectionPool.java | 9 +++++---- .../client/impl/conf/ClientConfigurationData.java | 8 +++++--- 5 files changed, 32 insertions(+), 10 deletions(-) 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();