From 9c795ebb596e4a62e22b8196321a125b0756c6b3 Mon Sep 17 00:00:00 2001 From: yfuruta Date: Fri, 19 Jul 2019 11:26:58 +0900 Subject: [PATCH] add timeout to internal rest api --- .../internal/NonPersistentTopicsImpl.java | 30 ++++++++++---- .../client/admin/internal/TopicsImpl.java | 40 ++++++++++++++----- 2 files changed, 53 insertions(+), 17 deletions(-) diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NonPersistentTopicsImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NonPersistentTopicsImpl.java index 8e81473bdcbed..d9b23c8f08157 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NonPersistentTopicsImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NonPersistentTopicsImpl.java @@ -23,6 +23,8 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import javax.ws.rs.client.Entity; import javax.ws.rs.client.InvocationCallback; @@ -52,12 +54,14 @@ public NonPersistentTopicsImpl(WebTarget web, Authentication auth, long readTime @Override public void createPartitionedTopic(String topic, int numPartitions) throws PulsarAdminException { try { - createPartitionedTopicAsync(topic, numPartitions).get(); + createPartitionedTopicAsync(topic, numPartitions).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -72,12 +76,14 @@ public CompletableFuture createPartitionedTopicAsync(String topic, int num @Override public PartitionedTopicMetadata getPartitionedTopicMetadata(String topic) throws PulsarAdminException { try { - return getPartitionedTopicMetadataAsync(topic).get(); + return getPartitionedTopicMetadataAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -105,12 +111,14 @@ public void failed(Throwable throwable) { @Override public NonPersistentTopicStats getStats(String topic) throws PulsarAdminException { try { - return getStatsAsync(topic).get(); + return getStatsAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -138,12 +146,14 @@ public void failed(Throwable throwable) { @Override public PersistentTopicInternalStats getInternalStats(String topic) throws PulsarAdminException { try { - return getInternalStatsAsync(topic).get(); + return getInternalStatsAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -171,12 +181,14 @@ public void failed(Throwable throwable) { @Override public void unload(String topic) throws PulsarAdminException { try { - unloadAsync(topic).get(); + unloadAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -190,12 +202,14 @@ public CompletableFuture unloadAsync(String topic) { @Override public List getListInBundle(String namespace, String bundleRange) throws PulsarAdminException { try { - return getListInBundleAsync(namespace, bundleRange).get(); + return getListInBundleAsync(namespace, bundleRange).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -221,12 +235,14 @@ public void failed(Throwable throwable) { @Override public List getList(String namespace) throws PulsarAdminException { try { - return getListAsync(namespace).get(); + return getListAsync(namespace).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java index 50f3c698b7501..79b9212918210 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java @@ -131,12 +131,14 @@ public List getPartitionedTopicList(String namespace) throws PulsarAdmin @Override public List getListInBundle(String namespace, String bundleRange) throws PulsarAdminException { try { - return getListInBundleAsync(namespace, bundleRange).get(); + return getListInBundleAsync(namespace, bundleRange).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -197,24 +199,28 @@ public void revokePermissions(String topic, String role) throws PulsarAdminExcep @Override public void createPartitionedTopic(String topic, int numPartitions) throws PulsarAdminException { try { - createPartitionedTopicAsync(topic, numPartitions).get(); + createPartitionedTopicAsync(topic, numPartitions).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @Override public void createNonPartitionedTopic(String topic) throws PulsarAdminException { try { - createNonPartitionedTopicAsync(topic).get(); + createNonPartitionedTopicAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -236,12 +242,14 @@ public CompletableFuture createPartitionedTopicAsync(String topic, int num @Override public void updatePartitionedTopic(String topic, int numPartitions) throws PulsarAdminException { try { - updatePartitionedTopicAsync(topic, numPartitions).get(); + updatePartitionedTopicAsync(topic, numPartitions).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -256,12 +264,14 @@ public CompletableFuture updatePartitionedTopicAsync(String topic, int num @Override public PartitionedTopicMetadata getPartitionedTopicMetadata(String topic) throws PulsarAdminException { try { - return getPartitionedTopicMetadataAsync(topic).get(); + return getPartitionedTopicMetadataAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -294,12 +304,14 @@ public void deletePartitionedTopic(String topic) throws PulsarAdminException { @Override public void deletePartitionedTopic(String topic, boolean force) throws PulsarAdminException { try { - deletePartitionedTopicAsync(topic, force).get(); + deletePartitionedTopicAsync(topic, force).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -319,12 +331,14 @@ public void delete(String topic) throws PulsarAdminException { @Override public void delete(String topic, boolean force) throws PulsarAdminException { try { - deleteAsync(topic, force).get(); + deleteAsync(topic, force).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -339,12 +353,14 @@ public CompletableFuture deleteAsync(String topic, boolean force) { @Override public void unload(String topic) throws PulsarAdminException { try { - unloadAsync(topic).get(); + unloadAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -358,12 +374,14 @@ public CompletableFuture unloadAsync(String topic) { @Override public List getSubscriptions(String topic) throws PulsarAdminException { try { - return getSubscriptionsAsync(topic).get(); + return getSubscriptionsAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } } @@ -595,12 +613,14 @@ public CompletableFuture deleteSubscriptionAsync(String topic, String subN @Override public void skipAllMessages(String topic, String subName) throws PulsarAdminException { try { - skipAllMessagesAsync(topic, subName).get(); + skipAllMessagesAsync(topic, subName).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); } catch (ExecutionException e) { throw (PulsarAdminException) e.getCause(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); } }