From e94923ad40899a07a166a8abd4d6c50850a3a2dc 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) # Conflicts: # pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java --- .../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 8f9b0ac55a92d..794e9df6b7d12 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 @@ -3544,7 +3544,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 011b47464ebfe..4872fbd2f0d16 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 @@ -2027,6 +2027,7 @@ public void addFailed(ManagedLedgerException exception, Object ctx) { @Test public void testCheckInactiveSubscriptions() throws Exception { + pulsarTestContext.getConfig().setAdditionalSystemCursorNames(Set.of("additionalSystemCursor")); PersistentTopic topic = new PersistentTopic(successTopicName, ledgerMock, brokerService); final var subscriptions = new ConcurrentHashMap(); @@ -2045,6 +2046,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); @@ -2071,6 +2077,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