From d09c0661b6e35277f2d8c182166e67d8f7a75455 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 16 Jan 2025 21:51:12 +0200 Subject: [PATCH 1/3] [fix][broker] Revert "[fix][broker] Cancel possible pending replay read in cancelPendingRead (#23384)" This reverts commit d2c91b1e1a8fc2fb233eb2856ddb6f53511ba201. --- .../persistent/PersistentDispatcherMultipleConsumers.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index b1cd186c31784..fa03a260e131e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -687,9 +687,8 @@ public synchronized CompletableFuture disconnectAllConsumers( @Override protected void cancelPendingRead() { - if ((havePendingRead || havePendingReplayRead) && cursor.cancelPendingReadRequest()) { + if (havePendingRead && cursor.cancelPendingReadRequest()) { havePendingRead = false; - havePendingReplayRead = false; } } From cd14e7b22704bbdd9f22aa0d3eebdff0e1e3d3d2 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 16 Jan 2025 22:10:09 +0200 Subject: [PATCH 2/3] Revert similar change in PersistentDispatcherMultipleConsumersClassic --- .../PersistentDispatcherMultipleConsumersClassic.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java index 6ab7acfa56da8..910491e60b2cf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java @@ -600,9 +600,8 @@ public synchronized CompletableFuture disconnectAllConsumers( @Override protected void cancelPendingRead() { - if ((havePendingRead || havePendingReplayRead) && cursor.cancelPendingReadRequest()) { + if (havePendingRead && cursor.cancelPendingReadRequest()) { havePendingRead = false; - havePendingReplayRead = false; } } From 41ad1ff9866d858a32a1bc3d0b58a7a78cfbe5a1 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 16 Jan 2025 22:10:37 +0200 Subject: [PATCH 3/3] Add javadoc to explain the purpose of "cancelPendingRead" --- .../broker/service/AbstractDispatcherMultipleConsumers.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java index e3c2cf40cf318..bec02e94c79ab 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java @@ -68,6 +68,10 @@ public SubType getType() { public abstract boolean isConsumerAvailable(Consumer consumer); + /** + * Cancel a possible pending read that is a Managed Cursor waiting to be notified for more entries. + * This won't cancel any other pending reads that are currently in progress. + */ protected void cancelPendingRead() {} /**