You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@pulsar.apache.org by mm...@apache.org on 2020/11/27 16:49:05 UTC

[pulsar] branch master updated: Add timeout to hasMessageAvailable to leader election process in Pulsar Functions (#8687)

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/pulsar.git


The following commit(s) were added to refs/heads/master by this push:
     new 9443ba2  Add timeout to hasMessageAvailable to leader election process in Pulsar Functions (#8687)
9443ba2 is described below

commit 9443ba2d0c63d8b170f2cd2390949a27baa171c2
Author: Boyang Jerry Peng <je...@gmail.com>
AuthorDate: Fri Nov 27 08:48:48 2020 -0800

    Add timeout to hasMessageAvailable to leader election process in Pulsar Functions (#8687)
    
    Co-authored-by: Jerry Peng <je...@splunk.com>
---
 .../org/apache/pulsar/functions/worker/FunctionAssignmentTailer.java    | 2 +-
 .../org/apache/pulsar/functions/worker/FunctionMetaDataTopicTailer.java | 2 +-
 2 files changed, 2 insertions(+), 2 deletions(-)

diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionAssignmentTailer.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionAssignmentTailer.java
index b76a290..48f087f 100644
--- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionAssignmentTailer.java
+++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionAssignmentTailer.java
@@ -145,7 +145,7 @@ public class FunctionAssignmentTailer implements AutoCloseable {
                 try {
                     Message<byte[]> msg = reader.readNext(1, TimeUnit.SECONDS);
                     if (msg == null) {
-                        if (exitOnEndOfTopic && !reader.hasMessageAvailable()) {
+                        if (exitOnEndOfTopic && !reader.hasMessageAvailableAsync().get(10, TimeUnit.SECONDS)) {
                             break;
                         }
                     } else {
diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionMetaDataTopicTailer.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionMetaDataTopicTailer.java
index 51b48df..934b589 100644
--- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionMetaDataTopicTailer.java
+++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/FunctionMetaDataTopicTailer.java
@@ -67,7 +67,7 @@ public class FunctionMetaDataTopicTailer
             try {
                 Message<byte[]> msg = reader.readNext(1, TimeUnit.SECONDS);
                 if (msg == null) {
-                    if (exitOnEndOfTopic && !reader.hasMessageAvailable()) {
+                    if (exitOnEndOfTopic && !reader.hasMessageAvailableAsync().get(10, TimeUnit.SECONDS)) {
                         break;
                     }
                 } else {