This is an automated email from the ASF dual-hosted git repository.

mmerli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git


The following commit(s) were added to refs/heads/master by this push:
     new d7ab96d  Move isConnected method into each interface class (#1780)
d7ab96d is described below

commit d7ab96d3ecb84e3e1e8ff83a1d6a9bceab38eaf6
Author: hrsakai <[email protected]>
AuthorDate: Wed May 16 01:59:24 2018 +0900

    Move isConnected method into each interface class (#1780)
---
 .../src/main/java/org/apache/pulsar/client/api/Consumer.java         | 5 +++++
 .../src/main/java/org/apache/pulsar/client/api/Producer.java         | 5 +++++
 pulsar-client/src/main/java/org/apache/pulsar/client/api/Reader.java | 5 +++++
 .../src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java    | 2 --
 .../src/main/java/org/apache/pulsar/client/impl/ProducerBase.java    | 2 --
 .../src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java      | 4 ++++
 6 files changed, 19 insertions(+), 4 deletions(-)

diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/Consumer.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/Consumer.java
index 876b6d2..69d885f 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/api/Consumer.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/api/Consumer.java
@@ -283,4 +283,9 @@ public interface Consumer<T> extends Closeable {
      * @return a future to track the completion of the seek operation
      */
     CompletableFuture<Void> seekAsync(MessageId messageId);
+
+    /**
+     * @return Whether the consumer is connected to the broker
+     */
+    boolean isConnected();
 }
diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/Producer.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/Producer.java
index 8f66f0e..a98073c 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/api/Producer.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/api/Producer.java
@@ -192,4 +192,9 @@ public interface Producer<T> extends Closeable {
      * @return a future that can used to track when the producer has been 
closed
      */
     CompletableFuture<Void> closeAsync();
+
+    /**
+     * @return Whether the producer is connected to the broker
+     */
+    boolean isConnected();
 }
diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/api/Reader.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/api/Reader.java
index 44c0d0a..cdf1132 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/api/Reader.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/api/Reader.java
@@ -73,4 +73,9 @@ public interface Reader<T> extends Closeable {
      * Asynchronously Check if there is message that has been published 
successfully to the broker in the topic.
      */
     CompletableFuture<Boolean> hasMessageAvailableAsync();
+
+    /**
+     * @return Whether the reader is connected to the broker
+     */
+    boolean isConnected();
 }
diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java
index cc718f3..016324e 100644
--- 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java
+++ 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java
@@ -317,8 +317,6 @@ public abstract class ConsumerBase<T> extends HandlerState 
implements Consumer<T
         return null;
     }
 
-    abstract public boolean isConnected();
-
     abstract public int getAvailablePermits();
 
     abstract public int numMessagesInQueue();
diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java
index d90f838..a75aa01 100644
--- 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java
+++ 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerBase.java
@@ -118,8 +118,6 @@ public abstract class ProducerBase<T> extends HandlerState 
implements Producer<T
     @Override
     abstract public CompletableFuture<Void> closeAsync();
 
-    abstract public boolean isConnected();
-
     @Override
     public String getTopic() {
         return topic;
diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java
index 9d0ed93..00f8af0 100644
--- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java
+++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ReaderImpl.java
@@ -148,4 +148,8 @@ public class ReaderImpl<T> implements Reader<T> {
         return consumer.hasMessageAvailableAsync();
     }
 
+    @Override
+    public boolean isConnected() {
+        return consumer.isConnected();
+    }
 }

-- 
To stop receiving notification emails like this one, please contact
[email protected].

Reply via email to