merlimat closed pull request #1780: Move isConnected method into each interface 
class
URL: https://github.com/apache/incubator-pulsar/pull/1780
 
 
   

This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:

As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):

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 876b6d21bc..69d885fb78 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 @@
      * @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 8f66f0ee18..a98073c141 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 @@
      * @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 44c0d0a553..cdf1132d6e 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 @@
      * 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 cc718f3d5f..016324e481 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 @@ protected SubType getSubType() {
         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 d90f838bab..a75aa01c38 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 void close() throws PulsarClientException {
     @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 9d0ed93d68..00f8af0b48 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 boolean hasMessageAvailable() throws 
PulsarClientException {
         return consumer.hasMessageAvailableAsync();
     }
 
+    @Override
+    public boolean isConnected() {
+        return consumer.isConnected();
+    }
 }


 

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
[email protected]


With regards,
Apache Git Services

Reply via email to