From adb8ccac18e1d99a34f6a4f27dc96bff40d34c8e Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Wed, 22 Feb 2017 09:20:01 -0800 Subject: [PATCH] If there is an error in reading from ZK, DiscoveryService should fail when reading partitions metadata --- .../service/BrokerDiscoveryProvider.java | 33 +++++++-------- .../service/DiscoveryServiceTest.java | 41 ++++++++++++++++--- .../pulsar/zookeeper/ZooKeeperCache.java | 8 +++- 3 files changed, 57 insertions(+), 25 deletions(-) diff --git a/pulsar-discovery-service/src/main/java/com/yahoo/pulsar/discovery/service/BrokerDiscoveryProvider.java b/pulsar-discovery-service/src/main/java/com/yahoo/pulsar/discovery/service/BrokerDiscoveryProvider.java index 89570f818a610..e465cb74dd809 100644 --- a/pulsar-discovery-service/src/main/java/com/yahoo/pulsar/discovery/service/BrokerDiscoveryProvider.java +++ b/pulsar-discovery-service/src/main/java/com/yahoo/pulsar/discovery/service/BrokerDiscoveryProvider.java @@ -39,7 +39,6 @@ import com.yahoo.pulsar.discovery.service.server.ServiceConfig; import com.yahoo.pulsar.discovery.service.web.ZookeeperCacheLoader; import com.yahoo.pulsar.zookeeper.GlobalZooKeeperCache; -import com.yahoo.pulsar.zookeeper.ZooKeeperCache.Deserializer; import com.yahoo.pulsar.zookeeper.ZooKeeperClientFactory; /** @@ -98,23 +97,21 @@ CompletableFuture getPartitionedTopicMetadata(Discover final String path = path(PARTITIONED_TOPIC_PATH_ZNODE, destination.getProperty(), destination.getCluster(), destination.getNamespacePortion(), "persistent", destination.getEncodedLocalName()); // gets the number of partitions from the zk cache - globalZkCache.getDataAsync(path, new Deserializer() { - @Override - public PartitionedTopicMetadata deserialize(String key, byte[] content) throws Exception { - return getThreadLocal().readValue(content, PartitionedTopicMetadata.class); - } - }).thenAccept(metadata -> { - // if the partitioned topic is not found in zk, then the topic - // is not partitioned - if (metadata.isPresent()) { - metadataFuture.complete(metadata.get()); - } else { - metadataFuture.complete(new PartitionedTopicMetadata()); - } - }).exceptionally(ex -> { - metadataFuture.complete(new PartitionedTopicMetadata()); - return null; - }); + globalZkCache + .getDataAsync(path, + (key, content) -> getThreadLocal().readValue(content, PartitionedTopicMetadata.class)) + .thenAccept(metadata -> { + // if the partitioned topic is not found in zk, then the topic + // is not partitioned + if (metadata.isPresent()) { + metadataFuture.complete(metadata.get()); + } else { + metadataFuture.complete(new PartitionedTopicMetadata()); + } + }).exceptionally(ex -> { + metadataFuture.completeExceptionally(ex); + return null; + }); } catch (Exception e) { metadataFuture.completeExceptionally(e); } diff --git a/pulsar-discovery-service/src/test/java/com/yahoo/pulsar/discovery/service/DiscoveryServiceTest.java b/pulsar-discovery-service/src/test/java/com/yahoo/pulsar/discovery/service/DiscoveryServiceTest.java index a2ca491dab1bb..aadc0f673161f 100644 --- a/pulsar-discovery-service/src/test/java/com/yahoo/pulsar/discovery/service/DiscoveryServiceTest.java +++ b/pulsar-discovery-service/src/test/java/com/yahoo/pulsar/discovery/service/DiscoveryServiceTest.java @@ -16,6 +16,7 @@ package com.yahoo.pulsar.discovery.service; import static com.yahoo.pulsar.discovery.service.web.ZookeeperCacheLoader.LOADBALANCE_BROKERS_ROOT; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; @@ -26,17 +27,23 @@ import java.net.URISyntaxException; import java.security.PrivateKey; import java.security.cert.X509Certificate; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import org.apache.bookkeeper.util.ZkUtils; import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.KeeperException.Code; +import org.apache.zookeeper.KeeperException.SessionExpiredException; import org.apache.zookeeper.ZooDefs; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; import com.yahoo.pulsar.common.api.Commands; +import com.yahoo.pulsar.common.naming.DestinationName; +import com.yahoo.pulsar.common.partition.PartitionedTopicMetadata; import com.yahoo.pulsar.common.policies.data.loadbalancer.LoadReport; import com.yahoo.pulsar.common.util.ObjectMapperFactory; import com.yahoo.pulsar.common.util.SecurityUtility; @@ -61,7 +68,7 @@ public class DiscoveryServiceTest extends BaseDiscoveryTestSetup { private final static String TLS_CLIENT_CERT_FILE_PATH = "./src/test/resources/certificate/client.crt"; private final static String TLS_CLIENT_KEY_FILE_PATH = "./src/test/resources/certificate/client.key"; - + @BeforeMethod private void init() throws Exception { super.setup(); @@ -88,6 +95,27 @@ public void testBrokerDiscoveryRoundRobin() throws Exception { } } + @Test + public void testGetPartitionsMetadata() throws Exception { + DestinationName topic1 = DestinationName.get("persistent://test/local/ns/my-topic-1"); + + PartitionedTopicMetadata m = service.getDiscoveryProvider().getPartitionedTopicMetadata(service, topic1, "role") + .get(); + assertEquals(m.partitions, 0); + + // Simulate ZK error + mockZookKeeper.failNow(Code.SESSIONEXPIRED); + DestinationName topic2 = DestinationName.get("persistent://test/local/ns/my-topic-2"); + CompletableFuture future = service.getDiscoveryProvider() + .getPartitionedTopicMetadata(service, topic2, "role"); + try { + future.get(); + fail("Partition metadata lookup should have failed"); + } catch (ExecutionException e) { + assertEquals(e.getCause().getClass(), SessionExpiredException.class); + } + } + /** * It verifies: client connects to Discovery-service and receives discovery response successfully. * @@ -137,18 +165,19 @@ public static NioEventLoopGroup connectToService(String serviceUrl, CountDownLat Bootstrap b = new Bootstrap(); b.group(workerGroup); b.channel(NioSocketChannel.class); - + b.handler(new ChannelInitializer() { @Override public void initChannel(SocketChannel ch) throws Exception { - if(tls) { + if (tls) { SslContextBuilder builder = SslContextBuilder.forClient(); builder.trustManager(InsecureTrustManagerFactory.INSTANCE); - X509Certificate[] certificates = SecurityUtility.loadCertificatesFromPemFile(TLS_CLIENT_CERT_FILE_PATH); + X509Certificate[] certificates = SecurityUtility + .loadCertificatesFromPemFile(TLS_CLIENT_CERT_FILE_PATH); PrivateKey privateKey = SecurityUtility.loadPrivateKeyFromPemFile(TLS_CLIENT_KEY_FILE_PATH); builder.keyManager(privateKey, (X509Certificate[]) certificates); SslContext sslCtx = builder.build(); - ch.pipeline().addLast("tls", sslCtx.newHandler(ch.alloc())); + ch.pipeline().addLast("tls", sslCtx.newHandler(ch.alloc())); } ch.pipeline().addLast(new ClientHandler(latch)); } @@ -156,7 +185,7 @@ public void initChannel(SocketChannel ch) throws Exception { URI uri = new URI(serviceUrl); InetSocketAddress serviceAddress = new InetSocketAddress(uri.getHost(), uri.getPort()); b.connect(serviceAddress).addListener((ChannelFuture future) -> { - if(!future.isSuccess()) { + if (!future.isSuccess()) { throw new IllegalStateException(future.cause()); } }); diff --git a/pulsar-zookeeper-utils/src/main/java/com/yahoo/pulsar/zookeeper/ZooKeeperCache.java b/pulsar-zookeeper-utils/src/main/java/com/yahoo/pulsar/zookeeper/ZooKeeperCache.java index a67b96dc367eb..f9e5b3099d52f 100644 --- a/pulsar-zookeeper-utils/src/main/java/com/yahoo/pulsar/zookeeper/ZooKeeperCache.java +++ b/pulsar-zookeeper-utils/src/main/java/com/yahoo/pulsar/zookeeper/ZooKeeperCache.java @@ -32,6 +32,7 @@ import org.apache.bookkeeper.util.SafeRunnable; import org.apache.zookeeper.KeeperException; import org.apache.zookeeper.KeeperException.Code; +import org.apache.zookeeper.KeeperException.NoNodeException; import org.apache.zookeeper.WatchedEvent; import org.apache.zookeeper.Watcher; import org.apache.zookeeper.ZooKeeper; @@ -209,7 +210,12 @@ public CompletableFuture> getDataAsync(final String path, final getDataAsync(path, this, deserializer).thenAccept(data -> { future.complete(data.map(e -> e.getKey())); }).exceptionally(ex -> { - future.complete(Optional.empty()); + if (ex.getCause() instanceof NoNodeException) { + future.complete(Optional.empty()); + } else { + future.completeExceptionally(ex.getCause()); + } + return null; }); return future;