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

nodece 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 803ff962018c [fix][broker] Read subscription properties directly from 
cursor (#26159)
803ff962018c is described below

commit 803ff962018c7e6229dd26af93cd44cfb5320459
Author: Zixuan Liu <[email protected]>
AuthorDate: Thu Jul 9 23:03:22 2026 +0800

    [fix][broker] Read subscription properties directly from cursor (#26159)
---
 .../service/persistent/PersistentSubscription.java | 21 ++++++++--------
 .../broker/service/persistent/PersistentTopic.java | 24 +++++++++++++++---
 .../persistent/PersistentSubscriptionTest.java     | 29 ++++++++++++++++++++++
 3 files changed, 60 insertions(+), 14 deletions(-)

diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
index e77ced8da168..fc8227cd153c 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
@@ -52,7 +52,6 @@ import 
org.apache.bookkeeper.mledger.ManagedLedgerException.ConcurrentFindCursor
 import 
org.apache.bookkeeper.mledger.ManagedLedgerException.InvalidCursorPositionException;
 import org.apache.bookkeeper.mledger.Position;
 import org.apache.bookkeeper.mledger.ScanOutcome;
-import org.apache.commons.collections4.MapUtils;
 import org.apache.commons.lang3.tuple.MutablePair;
 import org.apache.pulsar.broker.ServiceConfiguration;
 import org.apache.pulsar.broker.intercept.BrokerInterceptor;
@@ -133,7 +132,6 @@ public class PersistentSubscription extends 
AbstractSubscription {
     private volatile ReplicatedSubscriptionSnapshotCache 
replicatedSubscriptionSnapshotCache;
     @Getter
     private final PendingAckHandle pendingAckHandle;
-    private volatile Map<String, String> subscriptionProperties;
     private volatile CompletableFuture<Void> fenceFuture;
     private volatile CompletableFuture<Void> inProgressResetCursorFuture;
     private volatile Boolean replicatedControlled;
@@ -157,6 +155,14 @@ public class PersistentSubscription extends 
AbstractSubscription {
         this(topic, subscriptionName, cursor, replicated, 
Collections.emptyMap());
     }
 
+    /**
+     * Creates a persistent subscription.
+     *
+     * @deprecated use {@link #PersistentSubscription(PersistentTopic, String, 
ManagedCursor, Boolean)}
+     *             instead. The {@code subscriptionProperties} parameter is no 
longer read; the
+     *             cursor already carries all subscription properties.
+     */
+    @Deprecated
     public PersistentSubscription(PersistentTopic topic, String 
subscriptionName, ManagedCursor cursor,
                                   Boolean replicated, Map<String, String> 
subscriptionProperties) {
         this.topic = topic;
@@ -173,8 +179,6 @@ public class PersistentSubscription extends 
AbstractSubscription {
         if (replicated != null) {
             this.setReplicated(replicated);
         }
-        this.subscriptionProperties = MapUtils.isEmpty(subscriptionProperties)
-                ? Collections.emptyMap() : 
Collections.unmodifiableMap(subscriptionProperties);
         if (config.isTransactionCoordinatorEnabled()
                 && !isEventSystemTopic(TopicName.get(topicName))
                 && !ExtensibleLoadManagerImpl.isInternalTopic(topicName)) {
@@ -1516,7 +1520,7 @@ public class PersistentSubscription extends 
AbstractSubscription {
         subStats.msgExpired = expiryMonitor.getMessageExpiryCount();
         subStats.totalMsgExpired = expiryMonitor.getTotalMessageExpired();
         subStats.isReplicated = isReplicated();
-        subStats.subscriptionProperties = subscriptionProperties;
+        subStats.subscriptionProperties = getSubscriptionProperties();
         subStats.isDurable = cursor.isDurable();
         if (getType() == SubType.Key_Shared && dispatcher instanceof 
StickyKeyDispatcher) {
             StickyKeyDispatcher keySharedDispatcher = (StickyKeyDispatcher) 
dispatcher;
@@ -1657,7 +1661,7 @@ public class PersistentSubscription extends 
AbstractSubscription {
 
     @Override
     public Map<String, String> getSubscriptionProperties() {
-        return subscriptionProperties;
+        return cursor.getCursorProperties();
     }
 
     public Position getPositionInPendingAck(Position position) {
@@ -1671,10 +1675,7 @@ public class PersistentSubscription extends 
AbstractSubscription {
         } else {
             newSubscriptionProperties = 
Collections.unmodifiableMap(subscriptionProperties);
         }
-        return cursor.setCursorProperties(newSubscriptionProperties)
-                .thenRun(() -> {
-                    this.subscriptionProperties = newSubscriptionProperties;
-                });
+        return cursor.setCursorProperties(newSubscriptionProperties);
     }
     /**
      * Return a merged map that contains the cursor properties specified by 
used
diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
index 5eba52c8da50..7b5455e74583 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
@@ -1224,6 +1224,7 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
         synchronized (ledger) {
             // Create a new non-durable cursor only for the first consumer 
that connects
             PersistentSubscription subscription = 
subscriptions.get(subscriptionName);
+            CompletableFuture<Void> initPropertiesFuture = 
CompletableFuture.completedFuture(null);
 
             if (subscription == null) {
                 MessageIdImpl msgId = startMessageId != null ? (MessageIdImpl) 
startMessageId
@@ -1250,9 +1251,16 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
                     return FutureUtil.failedFuture(e);
                 }
 
-                subscription = new PersistentSubscription(this, 
subscriptionName, cursor, false,
-                        subscriptionProperties);
+                subscription = new PersistentSubscription(this, 
subscriptionName, cursor, false);
                 subscriptions.put(subscriptionName, subscription);
+
+                if (subscriptionProperties != null && 
!subscriptionProperties.isEmpty()) {
+                    // Trade-off: subscriptionProperties should be received by 
the cursor at creation time,
+                    // the way durable cursors take cursorProperties through 
ManagedLedger#asyncOpenCursor.
+                    // ManagedLedger#newNonDurableCursor has no equivalent 
parameter, so we seed the cursor
+                    // with a post-construction setCursorProperties call.
+                    initPropertiesFuture = 
cursor.setCursorProperties(subscriptionProperties);
+                }
             } else {
                 // if subscription exists, check if it's a durable subscription
                 if (subscription.getCursor() != null && 
subscription.getCursor().isDurable()) {
@@ -1261,11 +1269,19 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
                 }
             }
 
+            final PersistentSubscription finalSubscription = subscription;
             if (startMessageRollbackDurationSec > 0) {
-                resetSubscriptionCursor(subscription, subscriptionFuture, 
startMessageRollbackDurationSec);
+                initPropertiesFuture.whenComplete((__, ex) -> {
+                    if (ex != null) {
+                        subscriptionFuture.completeExceptionally(ex);
+                    } else {
+                        resetSubscriptionCursor(finalSubscription, 
subscriptionFuture,
+                                startMessageRollbackDurationSec);
+                    }
+                });
                 return subscriptionFuture;
             } else {
-                return CompletableFuture.completedFuture(subscription);
+                return initPropertiesFuture.thenApply(__ -> finalSubscription);
             }
         }
     }
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentSubscriptionTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentSubscriptionTest.java
index 0885a3ae43ef..a5437e8164cf 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentSubscriptionTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentSubscriptionTest.java
@@ -18,6 +18,7 @@
  */
 package org.apache.pulsar.broker.service.persistent;
 
+import static 
org.apache.bookkeeper.mledger.ManagedCursor.CURSOR_INTERNAL_PROPERTY_PREFIX;
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doAnswer;
@@ -34,6 +35,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import org.apache.bookkeeper.mledger.AsyncCallbacks;
@@ -48,6 +50,7 @@ import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl;
 import org.apache.commons.lang3.tuple.MutablePair;
 import org.apache.pulsar.broker.resources.NamespaceResources;
 import org.apache.pulsar.broker.service.Consumer;
+import org.apache.pulsar.broker.service.GetStatsOptions;
 import org.apache.pulsar.broker.testcontext.PulsarTestContext;
 import 
org.apache.pulsar.broker.transaction.buffer.impl.InMemTransactionBufferProvider;
 import org.apache.pulsar.broker.transaction.pendingack.PendingAckStore;
@@ -59,6 +62,7 @@ import org.apache.pulsar.common.api.proto.CommandAck.AckType;
 import org.apache.pulsar.common.api.proto.CommandSubscribe;
 import org.apache.pulsar.common.api.proto.TxnAction;
 import org.apache.pulsar.common.policies.data.Policies;
+import org.apache.pulsar.common.policies.data.stats.SubscriptionStatsImpl;
 import 
org.apache.pulsar.transaction.common.exception.TransactionConflictException;
 import org.awaitility.Awaitility;
 import org.testng.annotations.AfterMethod;
@@ -266,6 +270,31 @@ public class PersistentSubscriptionTest {
         assertThat(replicatedSubscriptionConfiguration).isEmpty();
     }
 
+    @Test
+    public void testGetSubscriptionPropertiesReflectsLiveCursorUpdates() 
throws Exception {
+        Map<String, String> backing = new ConcurrentHashMap<>();
+        doReturn(backing).when(cursorMock).getCursorProperties();
+        doAnswer(inv -> {
+            backing.put(inv.getArgument(0), inv.getArgument(1));
+            return CompletableFuture.completedFuture(null);
+        }).when(cursorMock).putCursorProperty(any(), any());
+        doReturn(false).when(cursorMock).isDurable();
+
+        
assertThat(persistentSubscription.getSubscriptionProperties()).isEmpty();
+
+        String bucketKey = CURSOR_INTERNAL_PROPERTY_PREFIX + 
"delayed.bucket_100_100";
+        persistentSubscription.getCursor().putCursorProperty(bucketKey, 
"42").get();
+
+        Map<String, String> live = 
persistentSubscription.getSubscriptionProperties();
+        assertThat(live).containsEntry(bucketKey, "42");
+
+        SubscriptionStatsImpl stats = persistentSubscription
+                .getStatsAsync(new GetStatsOptions(false, false, false, false, 
false)).get();
+        assertThat(stats.subscriptionProperties)
+                .isSameAs(live)
+                .containsEntry(bucketKey, "42");
+    }
+
     public static class CustomTransactionPendingAckStoreProvider implements 
TransactionPendingAckStoreProvider {
         @Override
         public CompletableFuture<PendingAckStore> 
newPendingAckStore(PersistentSubscription subscription) {

Reply via email to