From 1445fdba2d87eef872a3bae9123e90baa596e491 Mon Sep 17 00:00:00 2001 From: Yan Zhao Date: Thu, 16 Jul 2026 11:53:05 +0800 Subject: [PATCH] [improve][broker] Skip system cursor when check inactive cursor. (#26149) (cherry picked from commit fe61afda2d5bf4cb2de019cca47ef78cd1e350f1) --- .../pulsar/broker/service/persistent/PersistentTopic.java | 2 +- .../apache/pulsar/broker/service/PersistentTopicTest.java | 7 +++++++ 2 files changed, 8 insertions(+), 1 deletion(-) 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 fd45656fad44e..286b3146b0eb4 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 @@ -3524,7 +3524,7 @@ public void checkInactiveSubscriptions(long expirationTimeMillis) { subscriptions.forEach((subName, sub) -> { if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() || sub.isReplicated() - || isCompactionSubscription(subName)) { + || isSystemCursor(subName)) { return; } if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index 6ce87968155b8..62c30dd653560 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -2069,6 +2069,7 @@ public void addFailed(ManagedLedgerException exception, Object ctx) { @Test public void testCheckInactiveSubscriptions() throws Exception { pulsarTestContext.getConfig().setSubscriptionExpirationTimeMinutes(5); + pulsarTestContext.getConfig().setAdditionalSystemCursorNames(Set.of("additionalSystemCursor")); PersistentTopic topic = new PersistentTopic(successTopicName, ledgerMock, brokerService); final var subscriptions = new ConcurrentHashMap(); @@ -2087,6 +2088,11 @@ public void testCheckInactiveSubscriptions() throws Exception { spyWithClassAndConstructorArgsRecordingInvocations(PersistentSubscription.class, topic, "nonDeletableSubscription2", cursorMock, true); subscriptions.put(nonDeletableSubscription2.getName(), nonDeletableSubscription2); + // This subscription is an additional system cursor. + PersistentSubscription nonDeletableSubscription3 = + spyWithClassAndConstructorArgsRecordingInvocations(PersistentSubscription.class, topic, + "additionalSystemCursor", cursorMock, false); + subscriptions.put(nonDeletableSubscription3.getName(), nonDeletableSubscription3); Field field = topic.getClass().getDeclaredField("subscriptions"); field.setAccessible(true); @@ -2111,6 +2117,7 @@ public void testCheckInactiveSubscriptions() throws Exception { verify(nonDeletableSubscription1, times(0)).delete(); verify(deletableSubscription1, times(1)).delete(); verify(nonDeletableSubscription2, times(0)).delete(); + verify(nonDeletableSubscription3, times(0)).delete(); } @Test