Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -207,7 +207,7 @@

private final LoadingCache<String, SchemaInfoProvider> schemaProviderLoadingCache =
CacheBuilder.newBuilder().maximumSize(100000)
.expireAfterAccess(30, TimeUnit.MINUTES)

Check warning on line 210 in pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java

View workflow job for this annotation

GitHub Actions / Build and License check

[deprecation] expireAfterAccess(long,TimeUnit) in CacheBuilder has been deprecated

Check warning on line 210 in pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java

View workflow job for this annotation

GitHub Actions / Flaky tests suite

[deprecation] expireAfterAccess(long,TimeUnit) in CacheBuilder has been deprecated

Check warning on line 210 in pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java

View workflow job for this annotation

GitHub Actions / CI - Unit - Protobuf v3

[deprecation] expireAfterAccess(long,TimeUnit) in CacheBuilder has been deprecated

Check warning on line 210 in pulsar-client/src/main/java/org/apache/pulsar/client/impl/PulsarClientImpl.java

View workflow job for this annotation

GitHub Actions / Build Pulsar on MacOS

[deprecation] expireAfterAccess(long,TimeUnit) in CacheBuilder has been deprecated
.build(new CacheLoader<String, SchemaInfoProvider>() {

@Override
Expand Down Expand Up @@ -1496,7 +1496,9 @@
conf.getServiceUrlProvider().close();
}

if (addressResolver != null) {
// A DnsResolverGroupImpl hands the same resolver to every client on an event loop, so a resolver from a
// shared group is left to the group: closing it would break DNS for the clients that still use it.
if (addressResolver != null && dnsResolverGroupLocalInstance != null) {
Comment thread
lhotari marked this conversation as resolved.
try {
addressResolver.close();
} catch (Throwable t) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,9 @@
import io.netty.util.HashedWheelTimer;
import io.netty.util.concurrent.DefaultThreadFactory;
import java.lang.reflect.Field;
import java.net.DatagramPacket;
import java.net.DatagramSocket;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import java.util.ArrayList;
Expand All @@ -53,10 +56,12 @@
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit;
import java.util.regex.Pattern;
import lombok.Cleanup;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.PulsarClientSharedResources;
import org.apache.pulsar.client.api.ServiceUrlProvider;
import org.apache.pulsar.client.impl.conf.ClientConfigurationData;
import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData;
Expand Down Expand Up @@ -217,6 +222,37 @@ public void testShutdownContinuesWhenAddressResolverCloseFails() throws Exceptio
verify(timer).stop();
}

@Test(timeOut = 30_000)
public void testClosingClientLeavesSharedDnsResolverOpenForOtherClients() throws Exception {
// A DNS server that never answers: it only shows whether a resolver still sends queries
@Cleanup
DatagramSocket dnsServer = new DatagramSocket(0, InetAddress.getLoopbackAddress());
dnsServer.setSoTimeout(10_000);
@Cleanup
PulsarClientSharedResources sharedResources = PulsarClientSharedResources.builder()
// One event loop, so that both clients get the group's same resolver
.configureEventLoop(config -> config.numberOfThreads(1))
.configureDnsResolver(config -> config
.serverAddresses(List.of((InetSocketAddress) dnsServer.getLocalSocketAddress()))
.searchDomains(List.of())
.queryTimeoutMillis(TimeUnit.SECONDS.toMillis(20)))
.build();
PulsarClientImpl closedClient = (PulsarClientImpl) PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650").sharedResources(sharedResources).build();
@Cleanup
PulsarClientImpl client = (PulsarClientImpl) PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650").sharedResources(sharedResources).build();
assertSame(client.getAddressResolver(), closedClient.getAddressResolver());

closedClient.close();

// A closed resolver would fail the lookup without querying the DNS server
client.getAddressResolver().resolve(InetSocketAddress.createUnresolved("broker.pulsar.invalid", 6650));
DatagramPacket query = new DatagramPacket(new byte[512], 512);
dnsServer.receive(query);
assertTrue(query.getLength() > 0);
}

@Test
public void testInitializeWithTimer() throws PulsarClientException {
ClientConfigurationData conf = new ClientConfigurationData();
Expand Down
Loading