This is an automated email from the ASF dual-hosted git repository.

lhotari 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 e85b1f7d3b6 [fix][client] Fix UnAckedMessageRedeliveryTracker to skip 
cancelled timeouts (#26043)
e85b1f7d3b6 is described below

commit e85b1f7d3b69fb418afbf467f53ac58ce876d07c
Author: Dream95 <[email protected]>
AuthorDate: Fri Jul 3 21:35:56 2026 +0800

    [fix][client] Fix UnAckedMessageRedeliveryTracker to skip cancelled 
timeouts (#26043)
    
    Signed-off-by: Dream95 <[email protected]>
---
 .../pulsar/client/impl/UnAckedMessageRedeliveryTracker.java   | 11 +++++++++--
 1 file changed, 9 insertions(+), 2 deletions(-)

diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/UnAckedMessageRedeliveryTracker.java
 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/UnAckedMessageRedeliveryTracker.java
index 39fda87f024..7e425cdae89 100644
--- 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/UnAckedMessageRedeliveryTracker.java
+++ 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/UnAckedMessageRedeliveryTracker.java
@@ -57,6 +57,10 @@ public class UnAckedMessageRedeliveryTracker extends 
UnAckedMessageTracker {
         timeout = client.timer().newTimeout(new TimerTask() {
             @Override
             public void run(Timeout t) throws Exception {
+                if (t.isCancelled()) {
+                    return;
+                }
+
                 writeLock.lock();
                 try {
                     HashSet<UnackMessageIdWrapper> headPartition = 
redeliveryTimePartitions.removeFirst();
@@ -71,8 +75,11 @@ public class UnAckedMessageRedeliveryTracker extends 
UnAckedMessageTracker {
                     redeliveryTimePartitions.addLast(headPartition);
                     triggerRedelivery(consumerBase);
                 } finally {
-                    writeLock.unlock();
-                    timeout = client.timer().newTimeout(this, 
tickDurationInMs, TimeUnit.MILLISECONDS);
+                    try {
+                        timeout = client.timer().newTimeout(this, 
tickDurationInMs, TimeUnit.MILLISECONDS);
+                    } finally {
+                        writeLock.unlock();
+                    }
                 }
             }
         }, this.tickDurationInMs, TimeUnit.MILLISECONDS);

Reply via email to