Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions pulsar-client-api/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@
<artifactId>pulsar-transaction-common</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-resolver-dns</artifactId>
</dependency>
</dependencies>

<build>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -387,4 +388,12 @@ ClientBuilder authentication(String authPluginClassName, Map<String, String> 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);
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand All @@ -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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,17 +19,17 @@
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;
import org.apache.pulsar.client.api.Authentication;
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.
*/
Expand Down Expand Up @@ -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();

Expand Down