From 420d059621973cb8d178ade25c2beeb396399167 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 12 Aug 2023 14:40:43 +0800 Subject: [PATCH 1/5] fix compaction subscription delete by inactive subscription check --- .../service/persistent/PersistentTopic.java | 25 +++++---- .../broker/service/BrokerServiceTest.java | 52 +++++++++++++++++++ 2 files changed, 68 insertions(+), 9 deletions(-) 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 2b1700f36cb68..91fffb0f15190 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 @@ -2730,15 +2730,7 @@ public void checkInactiveSubscriptions() { final long expirationTimeMillis = TimeUnit.MINUTES .toMillis(nsExpirationTime == null ? defaultExpirationTime : nsExpirationTime); if (expirationTimeMillis > 0) { - subscriptions.forEach((subName, sub) -> { - if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() || sub.isReplicated()) { - return; - } - if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { - sub.delete().thenAccept(v -> log.info("[{}][{}] The subscription was deleted due to expiration " - + "with last active [{}]", topic, subName, sub.cursor.getLastActive())); - } - }); + checkInactiveSubscriptionsWithExpirationTime(expirationTimeMillis); } } catch (Exception e) { if (log.isDebugEnabled()) { @@ -2747,6 +2739,21 @@ public void checkInactiveSubscriptions() { } } + @VisibleForTesting + public void checkInactiveSubscriptionsWithExpirationTime(long expirationTimeMillis) { + subscriptions.forEach((subName, sub) -> { + if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() || sub.isReplicated() + || isCompactionSubscription(subName)) { + return; + } + + if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { + sub.delete().thenAccept(v -> log.info("[{}][{}] The subscription was deleted due to expiration " + + "with last active [{}]", topic, subName, sub.cursor.getLastActive())); + } + }); + } + @Override public void checkBackloggedCursors() { subscriptions.forEach((subName, subscription) -> { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 3d67a4bc840a5..53e96f9992a0b 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -77,6 +77,7 @@ import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitStateChannelImpl; import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.broker.service.BrokerServiceException.PersistenceException; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.stats.prometheus.PrometheusRawMetricsProvider; import org.apache.pulsar.client.admin.BrokerStats; @@ -109,6 +110,7 @@ import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.util.netty.EventLoopUtil; +import org.apache.pulsar.compaction.Compactor; import org.awaitility.Awaitility; import org.mockito.Mockito; import org.testng.Assert; @@ -1218,6 +1220,56 @@ public void testConcurrentLoadTopicExceedLimitShouldNotBeAutoCreated() throws Ex } } + @Test + public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() throws Exception { + String namespace = "prop/test"; + + // set up broker disable auto create and set concurrent load to 1 qps. + cleanup(); + conf.setBrokerServiceCompactionThresholdInBytes(8); + setup(); + + try { + admin.namespaces().createNamespace(namespace); + } catch (PulsarAdminException.ConflictException e) { + // Ok.. (if test fails intermittently and namespace is already created) + } + + String compactionInactiveTestTopic = "persistent://prop/test/testCompactionCursorShouldNotDelete"; + + admin.topics().createNonPartitionedTopic(compactionInactiveTestTopic); + + CompletableFuture> topicCf = + pulsar.getBrokerService().getTopic(compactionInactiveTestTopic, true); + + Optional topicOptional = topicCf.get(); + assertTrue(topicOptional.isPresent()); + + PersistentTopic topic = (PersistentTopic) topicOptional.get(); + + PersistentSubscription sub = (PersistentSubscription) topic.getSubscription(Compactor.COMPACTION_SUBSCRIPTION); + assertNotNull(sub); + + topic.checkCompaction(); + + Field currentCompaction = PersistentTopic.class.getDeclaredField("currentCompaction"); + currentCompaction.setAccessible(true); + CompletableFuture compactionFuture = (CompletableFuture)currentCompaction.get(topic); + + compactionFuture.get(); + + Thread.sleep(2000); + + // this operation is async delete + topic.checkInactiveSubscriptionsWithExpirationTime(1000); + + // wait for async delete finish + Thread.sleep(2000); + + assertNotNull(topic.getSubscription(Compactor.COMPACTION_SUBSCRIPTION)); + + } + /** * Verifies brokerService should not have deadlock and successfully remove topic from topicMap on topic-failure and * it should not introduce deadlock while performing it. From 7b543344573e16068a1ae37d1a9908e0f6053249 Mon Sep 17 00:00:00 2001 From: wangjinlong Date: Fri, 18 Aug 2023 15:38:33 +0800 Subject: [PATCH 2/5] 1. fix unit test and avoid add new public method --- .../service/persistent/PersistentTopic.java | 37 +++++++++++-------- .../broker/service/BrokerServiceTest.java | 17 ++++++--- 2 files changed, 33 insertions(+), 21 deletions(-) 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 91fffb0f15190..3dc87b287e06a 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 @@ -48,6 +48,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicLongFieldUpdater; import java.util.function.BiFunction; import java.util.stream.Collectors; import javax.annotation.Nonnull; @@ -247,6 +248,12 @@ protected TopicStatsHelper initialValue() { // Record the last time a data message (ie: not an internal Pulsar marker) is published on the topic private volatile long lastDataMessagePublishedTimestamp = 0; + + // Record the total expired subscription, if enable subscription expiration feature + private static final AtomicLongFieldUpdater EXPIRED_SUBSCRIPTION_NUMBERS_UPDATER = + AtomicLongFieldUpdater.newUpdater(PersistentTopic.class, "expiredSubscriptionNumbers"); + private volatile long expiredSubscriptionNumbers = 0L; + @Getter private final ExecutorService orderedExecutor; @@ -2730,7 +2737,16 @@ public void checkInactiveSubscriptions() { final long expirationTimeMillis = TimeUnit.MINUTES .toMillis(nsExpirationTime == null ? defaultExpirationTime : nsExpirationTime); if (expirationTimeMillis > 0) { - checkInactiveSubscriptionsWithExpirationTime(expirationTimeMillis); + subscriptions.forEach((subName, sub) -> { + if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() || sub.isReplicated()) { + return; + } + if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { + sub.delete().thenAccept(v -> log.info("[{}][{}] The subscription was deleted due to expiration " + + "with last active [{}]", topic, subName, sub.cursor.getLastActive())); + EXPIRED_SUBSCRIPTION_NUMBERS_UPDATER.incrementAndGet(this); + } + }); } } catch (Exception e) { if (log.isDebugEnabled()) { @@ -2739,21 +2755,6 @@ public void checkInactiveSubscriptions() { } } - @VisibleForTesting - public void checkInactiveSubscriptionsWithExpirationTime(long expirationTimeMillis) { - subscriptions.forEach((subName, sub) -> { - if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() || sub.isReplicated() - || isCompactionSubscription(subName)) { - return; - } - - if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { - sub.delete().thenAccept(v -> log.info("[{}][{}] The subscription was deleted due to expiration " - + "with last active [{}]", topic, subName, sub.cursor.getLastActive())); - } - }); - } - @Override public void checkBackloggedCursors() { subscriptions.forEach((subName, subscription) -> { @@ -3629,6 +3630,10 @@ public long getLastDataMessagePublishedTimestamp() { return lastDataMessagePublishedTimestamp; } + public long getExpiredSubscriptionNumbers() { + return EXPIRED_SUBSCRIPTION_NUMBERS_UPDATER.get(this); + } + public Optional getShadowSourceTopic() { return Optional.ofNullable(shadowSourceTopic); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 53e96f9992a0b..093167c10e3c0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -67,6 +67,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; +import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.http.HttpResponse; @@ -1258,14 +1259,20 @@ public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() thro compactionFuture.get(); - Thread.sleep(2000); + ManagedCursorImpl cursor = (ManagedCursorImpl) sub.getCursor(); - // this operation is async delete - topic.checkInactiveSubscriptionsWithExpirationTime(1000); + // make cursor last active time to very small to check if it will be deleted + Field cursorField = ManagedCursorImpl.class.getDeclaredField("lastActive"); + cursorField.setAccessible(true); + cursorField.set(cursor, 0); - // wait for async delete finish - Thread.sleep(2000); + // trigger inactive check. + topic.checkInactiveSubscriptions(); + // if subscription deleted the result should be zero + assertEquals(0, topic.getExpiredSubscriptionNumbers()); + + // check if the subscription is exist. assertNotNull(topic.getSubscription(Compactor.COMPACTION_SUBSCRIPTION)); } From 5793903da915d11bac50bc20aa7f285fa478f6fb Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 19 Aug 2023 11:28:55 +0800 Subject: [PATCH 3/5] fix unit test, fix check style --- .../pulsar/broker/service/persistent/PersistentTopic.java | 4 +++- .../org/apache/pulsar/broker/service/BrokerServiceTest.java | 4 +++- 2 files changed, 6 insertions(+), 2 deletions(-) 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 3dc87b287e06a..be8ef087963de 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 @@ -2738,7 +2738,9 @@ public void checkInactiveSubscriptions() { .toMillis(nsExpirationTime == null ? defaultExpirationTime : nsExpirationTime); if (expirationTimeMillis > 0) { subscriptions.forEach((subName, sub) -> { - if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() || sub.isReplicated()) { + if (sub.dispatcher != null && sub.dispatcher.isConsumerConnected() + || sub.isReplicated() + || isCompactionSubscription(subName)) { return; } if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 093167c10e3c0..3bd0553b95868 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -1236,6 +1236,8 @@ public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() thro // Ok.. (if test fails intermittently and namespace is already created) } + admin.namespaces().setSubscriptionExpirationTime(namespace, 1); + String compactionInactiveTestTopic = "persistent://prop/test/testCompactionCursorShouldNotDelete"; admin.topics().createNonPartitionedTopic(compactionInactiveTestTopic); @@ -1272,7 +1274,7 @@ public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() thro // if subscription deleted the result should be zero assertEquals(0, topic.getExpiredSubscriptionNumbers()); - // check if the subscription is exist. + // check if the subscription exist. assertNotNull(topic.getSubscription(Compactor.COMPACTION_SUBSCRIPTION)); } From 2272ce6463023e4e166a3bdac7b1513c3046a7e7 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 19 Aug 2023 11:30:50 +0800 Subject: [PATCH 4/5] fix comment --- .../org/apache/pulsar/broker/service/BrokerServiceTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 3bd0553b95868..5953fe0cd2456 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -1225,7 +1225,7 @@ public void testConcurrentLoadTopicExceedLimitShouldNotBeAutoCreated() throws Ex public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() throws Exception { String namespace = "prop/test"; - // set up broker disable auto create and set concurrent load to 1 qps. + // set up broker set compaction threshold. cleanup(); conf.setBrokerServiceCompactionThresholdInBytes(8); setup(); @@ -1236,6 +1236,7 @@ public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() thro // Ok.. (if test fails intermittently and namespace is already created) } + // set enable subscription expiration. admin.namespaces().setSubscriptionExpirationTime(namespace, 1); String compactionInactiveTestTopic = "persistent://prop/test/testCompactionCursorShouldNotDelete"; From 128ff3fdb98b2e7338bb6822705bb773baefb4b4 Mon Sep 17 00:00:00 2001 From: WJL3333 Date: Sat, 19 Aug 2023 11:50:17 +0800 Subject: [PATCH 5/5] fix comment --- .../broker/service/persistent/PersistentTopic.java | 12 ------------ .../pulsar/broker/service/BrokerServiceTest.java | 14 +++++++++----- 2 files changed, 9 insertions(+), 17 deletions(-) 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 be8ef087963de..410e67c6858f0 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 @@ -48,7 +48,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; -import java.util.concurrent.atomic.AtomicLongFieldUpdater; import java.util.function.BiFunction; import java.util.stream.Collectors; import javax.annotation.Nonnull; @@ -248,12 +247,6 @@ protected TopicStatsHelper initialValue() { // Record the last time a data message (ie: not an internal Pulsar marker) is published on the topic private volatile long lastDataMessagePublishedTimestamp = 0; - - // Record the total expired subscription, if enable subscription expiration feature - private static final AtomicLongFieldUpdater EXPIRED_SUBSCRIPTION_NUMBERS_UPDATER = - AtomicLongFieldUpdater.newUpdater(PersistentTopic.class, "expiredSubscriptionNumbers"); - private volatile long expiredSubscriptionNumbers = 0L; - @Getter private final ExecutorService orderedExecutor; @@ -2746,7 +2739,6 @@ public void checkInactiveSubscriptions() { if (System.currentTimeMillis() - sub.cursor.getLastActive() > expirationTimeMillis) { sub.delete().thenAccept(v -> log.info("[{}][{}] The subscription was deleted due to expiration " + "with last active [{}]", topic, subName, sub.cursor.getLastActive())); - EXPIRED_SUBSCRIPTION_NUMBERS_UPDATER.incrementAndGet(this); } }); } @@ -3632,10 +3624,6 @@ public long getLastDataMessagePublishedTimestamp() { return lastDataMessagePublishedTimestamp; } - public long getExpiredSubscriptionNumbers() { - return EXPIRED_SUBSCRIPTION_NUMBERS_UPDATER.get(this); - } - public Optional getShadowSourceTopic() { return Optional.ofNullable(shadowSourceTopic); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java index 5953fe0cd2456..6497daa81dcfa 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceTest.java @@ -1265,15 +1265,19 @@ public void testCheckInactiveSubscriptionsShouldNotDeleteCompactionCursor() thro ManagedCursorImpl cursor = (ManagedCursorImpl) sub.getCursor(); // make cursor last active time to very small to check if it will be deleted - Field cursorField = ManagedCursorImpl.class.getDeclaredField("lastActive"); - cursorField.setAccessible(true); - cursorField.set(cursor, 0); + Field cursorLastActiveField = ManagedCursorImpl.class.getDeclaredField("lastActive"); + cursorLastActiveField.setAccessible(true); + cursorLastActiveField.set(cursor, 0); + + // replace origin object. so we can check if subscription is deleted. + PersistentSubscription spySubscription = Mockito.spy(sub); + topic.getSubscriptions().put(Compactor.COMPACTION_SUBSCRIPTION, spySubscription); // trigger inactive check. topic.checkInactiveSubscriptions(); - // if subscription deleted the result should be zero - assertEquals(0, topic.getExpiredSubscriptionNumbers()); + // Compaction subscription should not call delete method. + Mockito.verify(spySubscription, Mockito.never()).delete(); // check if the subscription exist. assertNotNull(topic.getSubscription(Compactor.COMPACTION_SUBSCRIPTION));