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].