From 96ca94525e6218774f50ac553162c86bee5e2c51 Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Sat, 8 Jan 2022 23:34:29 +0800 Subject: [PATCH] [Broker] Add operation timeout to metadata store Signed-off-by: Zixuan Liu --- .../broker/resources/PulsarResources.java | 19 +++++--- .../apache/pulsar/broker/PulsarService.java | 2 + .../proxy/ProxyAuthenticationTest.java | 3 +- .../proxy/ProxyAuthorizationTest.java | 3 +- .../proxy/ProxyConfigurationTest.java | 3 +- .../proxy/ProxyPublishConsumeTest.java | 3 +- .../proxy/ProxyPublishConsumeTlsTest.java | 3 +- .../ProxyPublishConsumeWithoutZKTest.java | 2 +- .../proxy/v1/V1_ProxyAuthenticationTest.java | 3 +- .../pulsar/functions/worker/Worker.java | 3 +- .../metadata/api/MetadataStoreConfig.java | 6 +++ .../metadata/impl/AbstractMetadataStore.java | 8 ++++ .../pulsar/metadata/impl/ZKMetadataStore.java | 22 +++++++--- .../AbstractBatchedMetadataStore.java | 5 ++- .../metadata/BaseMetadataStoreTest.java | 3 ++ .../metadata/impl/ZKMetadataStoreTest.java | 44 +++++++++++++++++++ .../pulsar/websocket/WebSocketService.java | 8 ++-- 17 files changed, 114 insertions(+), 26 deletions(-) create mode 100644 pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/ZKMetadataStoreTest.java diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java index 7a8e00a289ec3..41cfd86c8195a 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/resources/PulsarResources.java @@ -19,14 +19,12 @@ package org.apache.pulsar.broker.resources; import java.util.Optional; - +import lombok.Getter; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.MetadataStoreConfig; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; -import lombok.Getter; - public class PulsarResources { public static final int DEFAULT_OPERATION_TIMEOUT_SEC = 30; @@ -91,7 +89,18 @@ public PulsarResources(MetadataStore localMetadataStore, MetadataStore configura public static MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs) throws MetadataStoreException { - return MetadataStoreExtended.create(serverUrls, MetadataStoreConfig.builder() - .sessionTimeoutMillis(sessionTimeoutMs).allowReadOnlyOperations(false).build()); + return createMetadataStore(serverUrls, sessionTimeoutMs, 0); + } + + public static MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs, + int operationTimeoutSeconds) + throws MetadataStoreException { + MetadataStoreConfig.MetadataStoreConfigBuilder configBuilder = MetadataStoreConfig.builder() + .sessionTimeoutMillis(sessionTimeoutMs).allowReadOnlyOperations(false); + if (operationTimeoutSeconds > 0) { + configBuilder.operationTimeoutSeconds(operationTimeoutSeconds); + } + + return MetadataStoreExtended.create(serverUrls, configBuilder.build()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 666673e8166b1..7830ba06c574c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -329,6 +329,7 @@ public MetadataStore createConfigurationMetadataStore() throws MetadataStoreExce .batchingMaxDelayMillis(config.getMetadataStoreBatchingMaxDelayMillis()) .batchingMaxOperations(config.getMetadataStoreBatchingMaxOperations()) .batchingMaxSizeKb(config.getMetadataStoreBatchingMaxSizeKb()) + .operationTimeoutSeconds(config.getZooKeeperOperationTimeoutSeconds()) .build()); } @@ -913,6 +914,7 @@ public MetadataStoreExtended createLocalMetadataStore() throws MetadataStoreExce .batchingMaxDelayMillis(config.getMetadataStoreBatchingMaxDelayMillis()) .batchingMaxOperations(config.getMetadataStoreBatchingMaxOperations()) .batchingMaxSizeKb(config.getMetadataStoreBatchingMaxSizeKb()) + .operationTimeoutSeconds(config.getZooKeeperOperationTimeoutSeconds()) .build()); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java index 764d4a5b7915a..35fa98a02b8a3 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthenticationTest.java @@ -83,7 +83,8 @@ public void setup() throws Exception { } service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createMetadataStore(anyString(), anyInt(), anyInt()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java index 15c3e50e185b8..9b0852c570ffc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyAuthorizationTest.java @@ -65,7 +65,8 @@ protected void setup() throws Exception { config.setWebServicePort(Optional.of(0)); config.setConfigurationStoreServers(GLOBAL_DUMMY_VALUE); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createMetadataStore(anyString(), anyInt(), anyInt()); service.start(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java index aeeab2afddcce..601409492d15e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyConfigurationTest.java @@ -64,7 +64,8 @@ public void configTest(int numIoThreads, int connectionsPerBroker) throws Except config.setWebSocketNumIoThreads(numIoThreads); config.setWebSocketConnectionsPerBroker(connectionsPerBroker); WebSocketService service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createMetadataStore(anyString(), anyInt(), anyInt()); service.start(); PulsarClientImpl client = (PulsarClientImpl) service.getPulsarClient(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java index 177f6f8795a0f..313b155d99de9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTest.java @@ -101,7 +101,8 @@ public void setup() throws Exception { config.setClusterName("test"); config.setConfigurationStoreServers(GLOBAL_DUMMY_VALUE); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createMetadataStore(anyString(), anyInt(), anyInt()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java index b8702522fce13..5b4f68d15c59d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeTlsTest.java @@ -75,7 +75,8 @@ public void setup() throws Exception { config.setBrokerClientAuthenticationPlugin(AuthenticationTls.class.getName()); config.setConfigurationStoreServers(GLOBAL_DUMMY_VALUE); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createMetadataStore(anyString(), anyInt(), anyInt()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java index 1fb12645e5e34..5dc81189ad5c5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/ProxyPublishConsumeWithoutZKTest.java @@ -62,7 +62,7 @@ public void setup() throws Exception { config.setServiceUrl(pulsar.getSafeWebServiceAddress()); config.setServiceUrlTls(pulsar.getWebServiceAddressTls()); service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeper)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeper)).when(service).createMetadataStore(anyString(), anyInt(), anyInt()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java index 842acede2492c..e2420af9fb9f3 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/websocket/proxy/v1/V1_ProxyAuthenticationTest.java @@ -85,7 +85,8 @@ public void setup() throws Exception { } service = spyWithClassAndConstructorArgs(WebSocketService.class, config); - doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service).createMetadataStore(anyString(), anyInt()); + doReturn(new ZKMetadataStore(mockZooKeeperGlobal)).when(service) + .createMetadataStore(anyString(), anyInt(), anyInt()); proxyServer = new ProxyServer(config); WebSocketServiceStarter.start(proxyServer, service); log.info("Proxy Server Started"); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java index 929e62e862621..f30d607bcf0d7 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/Worker.java @@ -74,7 +74,8 @@ private AuthorizationService getAuthorizationService() throws PulsarServerExcept log.info("starting configuration cache service"); try { configMetadataStore = PulsarResources.createMetadataStore(workerConfig.getConfigurationStoreServers(), - (int) workerConfig.getZooKeeperSessionTimeoutMillis()); + (int) workerConfig.getZooKeeperSessionTimeoutMillis(), + workerConfig.getZooKeeperOperationTimeoutSeconds()); } catch (IOException e) { throw new PulsarServerException(e); } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java index 651b32f2464be..34cf945d6a856 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreConfig.java @@ -69,4 +69,10 @@ public class MetadataStoreConfig { */ @Builder.Default private final int batchingMaxSizeKb = 128; + + /** + * The operation timeout in seconds. + */ + @Builder.Default + private final int operationTimeoutSeconds = 30; } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java index e8230e0113ffe..c3e227a206dde 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java @@ -44,6 +44,7 @@ import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataCache; import org.apache.pulsar.metadata.api.MetadataSerde; +import org.apache.pulsar.metadata.api.MetadataStoreConfig; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.Notification; import org.apache.pulsar.metadata.api.NotificationType; @@ -64,6 +65,8 @@ public abstract class AbstractMetadataStore implements MetadataStoreExtended, Co private final AsyncLoadingCache> childrenCache; private final AsyncLoadingCache existsCache; private final CopyOnWriteArrayList> metadataCaches = new CopyOnWriteArrayList<>(); + @Getter + private final MetadataStoreConfig metadataStoreConfig; // We don't strictly need to use 'volatile' here because we don't need the precise consistent semantic. Instead, // we want to avoid the overhead of 'volatile'. @@ -75,6 +78,11 @@ public abstract class AbstractMetadataStore implements MetadataStoreExtended, Co protected abstract CompletableFuture existsFromStore(String path); protected AbstractMetadataStore() { + this(null); + } + + protected AbstractMetadataStore(MetadataStoreConfig metadataStoreConfig) { + this.metadataStoreConfig = metadataStoreConfig; this.executor = Executors .newSingleThreadScheduledExecutor(new DefaultThreadFactory("metadata-store")); registerListener(this); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java index cf0b7c3d049ed..dc11c32176efb 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java @@ -27,6 +27,8 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; import lombok.SneakyThrows; @@ -69,7 +71,6 @@ public class ZKMetadataStore extends AbstractBatchedMetadataStore implements MetadataStoreExtended, MetadataStoreLifecycle { private final String metadataURL; - private final MetadataStoreConfig metadataStoreConfig; private final boolean isZkManaged; private final ZooKeeper zkc; private Optional sessionWatcher; @@ -80,7 +81,6 @@ public ZKMetadataStore(String metadataURL, MetadataStoreConfig metadataStoreConf try { this.metadataURL = metadataURL; - this.metadataStoreConfig = metadataStoreConfig; isZkManaged = true; zkc = PulsarZooKeeperClient.newBuilder().connectString(metadataURL) .connectRetryPolicy(new BoundExponentialBackoffRetryPolicy(100, 60_000, Integer.MAX_VALUE)) @@ -109,7 +109,6 @@ public ZKMetadataStore(ZooKeeper zkc) { super(MetadataStoreConfig.builder().build()); this.metadataURL = null; - this.metadataStoreConfig = null; this.isZkManaged = false; this.zkc = zkc; this.sessionWatcher = Optional.of(new ZKSessionWatcher(zkc, this::receivedSessionEvent)); @@ -145,7 +144,10 @@ protected void receivedSessionEvent(SessionEvent event) { @Override protected void batchOperation(List ops) { try { + AtomicBoolean callback = new AtomicBoolean(false); + zkc.multi(ops.stream().map(this::convertOp).collect(Collectors.toList()), (rc, path, ctx, results) -> { + callback.set(true); if (results == null) { Code code = Code.get(rc); if (code == Code.CONNECTIONLOSS) { @@ -186,6 +188,12 @@ protected void batchOperation(List ops) { } } }, null); + + executor.schedule(() -> { + if (!callback.get()) { + ops.forEach(n -> n.getFuture().completeExceptionally(new TimeoutException())); + } + }, getMetadataStoreConfig().getOperationTimeoutSeconds(), TimeUnit.SECONDS); } catch (Throwable t) { ops.forEach(o -> o.getFuture().completeExceptionally(new MetadataStoreException(t))); } @@ -501,7 +509,7 @@ public CompletableFuture initializeCluster() { if (this.metadataURL == null) { return FutureUtil.failedFuture(new MetadataStoreException("metadataURL is not set")); } - if (this.metadataStoreConfig == null) { + if (this.getMetadataStoreConfig() == null) { return FutureUtil.failedFuture(new MetadataStoreException("metadataStoreConfig is not set")); } int chrootIndex = metadataURL.indexOf("/"); @@ -510,10 +518,10 @@ public CompletableFuture initializeCluster() { String zkConnectForChrootCreation = metadataURL.substring(0, chrootIndex); try (ZooKeeper chrootZk = PulsarZooKeeperClient.newBuilder() .connectString(zkConnectForChrootCreation) - .sessionTimeoutMs(metadataStoreConfig.getSessionTimeoutMillis()) + .sessionTimeoutMs(getMetadataStoreConfig().getSessionTimeoutMillis()) .connectRetryPolicy( - new BoundExponentialBackoffRetryPolicy(metadataStoreConfig.getSessionTimeoutMillis(), - metadataStoreConfig.getSessionTimeoutMillis(), 0)) + new BoundExponentialBackoffRetryPolicy(getMetadataStoreConfig().getSessionTimeoutMillis(), + getMetadataStoreConfig().getSessionTimeoutMillis(), 0)) .build()) { if (chrootZk.exists(chrootPath, false) == null) { createFullPathOptimistic(chrootZk, chrootPath, new byte[0], CreateMode.PERSISTENT); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index 5d78282abcac3..c24d93bf780b8 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -27,6 +27,7 @@ import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; +import lombok.NonNull; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataStoreConfig; @@ -50,8 +51,8 @@ public abstract class AbstractBatchedMetadataStore extends AbstractMetadataStore private final int maxOperations; private final int maxSize; - protected AbstractBatchedMetadataStore(MetadataStoreConfig conf) { - super(); + protected AbstractBatchedMetadataStore(@NonNull MetadataStoreConfig conf) { + super(conf); this.enabled = conf.isBatchingEnabled(); this.maxDelayMillis = conf.getBatchingMaxDelayMillis(); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BaseMetadataStoreTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BaseMetadataStoreTest.java index fcb6b77efe5ac..69185a932ce1a 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BaseMetadataStoreTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BaseMetadataStoreTest.java @@ -24,6 +24,7 @@ import java.util.concurrent.CompletionException; import java.util.function.Supplier; import org.apache.pulsar.tests.TestRetrySupport; +import org.apache.zookeeper.ZooKeeper; import org.assertj.core.util.Files; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; @@ -31,12 +32,14 @@ public abstract class BaseMetadataStoreTest extends TestRetrySupport { protected TestZKServer zks; + protected ZooKeeper zkc; @BeforeClass(alwaysRun = true) @Override public final void setup() throws Exception { incrementSetupNumber(); zks = new TestZKServer(); + zkc = new ZooKeeper(zks.getConnectionString(), 30_000, null); } @AfterClass(alwaysRun = true) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/ZKMetadataStoreTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/ZKMetadataStoreTest.java new file mode 100644 index 0000000000000..5213218f7ca5e --- /dev/null +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/ZKMetadataStoreTest.java @@ -0,0 +1,44 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.metadata.impl; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.spy; + +import java.util.concurrent.CompletionException; +import java.util.concurrent.TimeoutException; +import org.apache.pulsar.metadata.BaseMetadataStoreTest; +import org.apache.pulsar.metadata.api.MetadataStore; +import org.apache.zookeeper.ZooKeeper; +import org.testng.annotations.Test; + +public class ZKMetadataStoreTest extends BaseMetadataStoreTest { + @Test + public void testOperationTimeout() { + ZooKeeper zooKeeper = spy(zkc); + doNothing().when(zooKeeper).multi(any(), any(), any()); + + MetadataStore zkMetadataStore = new ZKMetadataStore(zooKeeper); + CompletionException ex = assertThrows(CompletionException.class, () -> zkMetadataStore.get("/").join()); + assertEquals(TimeoutException.class, ex.getCause().getClass()); + } +} diff --git a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java index 5abaf52bc73b8..ae42a02e7a9ac 100644 --- a/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java +++ b/pulsar-websocket/src/main/java/org/apache/pulsar/websocket/WebSocketService.java @@ -93,7 +93,7 @@ public void start() throws PulsarServerException, PulsarClientException, Malform if (isNotBlank(config.getConfigurationMetadataStoreUrl())) { try { configMetadataStore = createMetadataStore(config.getConfigurationMetadataStoreUrl(), - (int) config.getZooKeeperSessionTimeoutMillis()); + (int) config.getZooKeeperSessionTimeoutMillis(), config.getZooKeeperOperationTimeoutSeconds()); } catch (MetadataStoreException e) { throw new PulsarServerException(e); } @@ -113,9 +113,9 @@ public void start() throws PulsarServerException, PulsarClientException, Malform log.info("Pulsar WebSocket Service started"); } - public MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs) - throws MetadataStoreException { - return PulsarResources.createMetadataStore(serverUrls, sessionTimeoutMs); + public MetadataStoreExtended createMetadataStore(String serverUrls, int sessionTimeoutMs, + int operationTimeoutSeconds) throws MetadataStoreException { + return PulsarResources.createMetadataStore(serverUrls, sessionTimeoutMs, operationTimeoutSeconds); } @Override