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