From 007c0ef27903bcb0f56065a6500138bca6766cf6 Mon Sep 17 00:00:00 2001 From: Michael Darakananda Date: Mon, 5 Jun 2017 15:27:23 +1000 Subject: [PATCH 1/2] pubsub: eventually stop polling after termination Update #2103. This does not completely fix the problem. If an RPC is started before the shutdown, the RPC won't be canceled. Depending on how the channels and executors are set up, messages in the said RPCs might be lease-extended. With this PR, at least we're not polling forever anymore. --- .../cloud/pubsub/spi/v1/PollingSubscriberConnection.java | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java index da9b64282ac9..94c1df2035ba 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java @@ -127,6 +127,9 @@ protected void doStop() { } private void pullMessages(final Duration backoff) { + if (!isAlive()) { + return; + } ListenableFuture pullResult = stub.withDeadlineAfter(DEFAULT_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS) .pull( @@ -201,7 +204,9 @@ public void run() { } private boolean isAlive() { - return state() == State.RUNNING || state() == State.STARTING; + // Read state only once. Because of threading, different calss can give different results. + State state = state(); + return state == State.RUNNING || state == State.STARTING; } @Override From dbcb1ef6dc59b39c4c512c4132ac6e897f16b720 Mon Sep 17 00:00:00 2001 From: Michael Darakananda Date: Tue, 6 Jun 2017 12:19:03 +1000 Subject: [PATCH 2/2] pr comment --- .../google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java index 94c1df2035ba..f727972a2b7a 100644 --- a/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java +++ b/google-cloud-pubsub/src/main/java/com/google/cloud/pubsub/spi/v1/PollingSubscriberConnection.java @@ -204,7 +204,7 @@ public void run() { } private boolean isAlive() { - // Read state only once. Because of threading, different calss can give different results. + // Read state only once. Because of threading, different calls can give different results. State state = state(); return state == State.RUNNING || state == State.STARTING; }