diff --git a/conf/broker.conf b/conf/broker.conf index 16688069caa13..5fe9e8dccdb9c 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -452,6 +452,10 @@ systemTopicEnabled=false # Please enable the system topic first. topicLevelPoliciesEnabled=false +# If a topic remains fenced for this number of seconds, it will be closed forcefully. +# If it is set to 0 or a negative number, the fenced topic will not be closed. +topicFencingTimeoutSeconds=0 + ### --- Authentication --- ### # Role names that are treated as "proxy roles". If the broker sees a request with #role as proxyRoles - it will demand to see a valid original principal. diff --git a/conf/standalone.conf b/conf/standalone.conf index f3a5b4b3ba589..4cea72d82c128 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -365,6 +365,10 @@ systemTopicEnabled=false # Please enable the system topic first. topicLevelPoliciesEnabled=false +# If a topic remains fenced for this number of seconds, it will be closed forcefully. +# If it is set to 0 or a negative number, the fenced topic will not be closed. +topicFencingTimeoutSeconds=0 + ### --- Authentication --- ### # Role names that are treated as "proxy roles". If the broker sees a request with #role as proxyRoles - it will demand to see a valid original principal. diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 958dfea97d48e..7be222734c5b2 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -827,6 +827,12 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private String zookeeperSessionExpiredPolicy = "shutdown"; + @FieldContext( + category = CATEGORY_SERVER, + doc = "If a topic remains fenced for this number of seconds, it will be closed forcefully.\n" + + " If it is set to 0 or a negative number, the fenced topic will not be closed." + ) + private int topicFencingTimeoutSeconds = 0; /**** --- Messaging Protocols --- ****/ 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 ec7215e56d3a0..ea877dcb13a1b 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 @@ -40,6 +40,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; @@ -194,6 +195,8 @@ protected TopicStatsHelper initialValue() { private volatile int maxUnackedMessagesOnSubscription = -1; private volatile boolean isClosingOrDeleting = false; + private ScheduledFuture fencedTopicMonitoringTask = null; + private static class TopicStatsHelper { public double averageMsgSize; public double aggMsgRateIn; @@ -396,7 +399,7 @@ private void decrementPendingWriteOpsAndCheck() { // signal to managed ledger that we are ready to resume by creating a new ledger ledger.readyToCreateNewLedger(); - isFenced = false; + unfence(); } } @@ -423,7 +426,7 @@ public synchronized void addFailed(ManagedLedgerException exception, Object ctx) } else { // fence topic when failed to write a message to BK - isFenced = true; + fence(); // close all producers List> futures = Lists.newArrayList(); producers.values().forEach(producer -> futures.add(producer.disconnect())); @@ -2384,6 +2387,40 @@ public boolean isSystemTopic() { return false; } + private synchronized void fence() { + isFenced = true; + ScheduledFuture monitoringTask = this.fencedTopicMonitoringTask; + if (monitoringTask == null || monitoringTask.isDone()) { + final int timeout = brokerService.pulsar().getConfiguration().getTopicFencingTimeoutSeconds(); + if (timeout > 0) { + this.fencedTopicMonitoringTask = brokerService.executor().schedule(this::closeFencedTopicForcefully, + timeout, TimeUnit.SECONDS); + } + } + } + + private synchronized void unfence() { + isFenced = false; + ScheduledFuture monitoringTask = this.fencedTopicMonitoringTask; + if (monitoringTask != null && !monitoringTask.isDone()) { + monitoringTask.cancel(false); + } + } + + private void closeFencedTopicForcefully() { + if (isFenced) { + final int timeout = brokerService.pulsar().getConfiguration().getTopicFencingTimeoutSeconds(); + if (isClosingOrDeleting) { + log.warn("[{}] Topic remained fenced for {} seconds and is already closed (pendingWriteOps: {})", topic, + timeout, pendingWriteOps.get()); + } else { + log.error("[{}] Topic remained fenced for {} seconds, so close it (pendingWriteOps: {})", topic, + timeout, pendingWriteOps.get()); + close(); + } + } + } + private void fenceTopicToCloseOrDelete() { isClosingOrDeleting = true; isFenced = true; 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 74a5cb57388a7..161731f3948fa 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 @@ -59,6 +59,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.ForkJoinPool; import java.util.concurrent.Future; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; @@ -1749,6 +1750,43 @@ public void testCheckInactiveSubscriptions() throws Exception { verify(nonDeletableSubscription2, times(0)).delete(); } + @Test + public void testTopicFencingTimeout() throws Exception { + ServiceConfiguration svcConfig = spy(new ServiceConfiguration()); + doReturn(svcConfig).when(pulsar).getConfiguration(); + PersistentTopic topic = new PersistentTopic(successTopicName, ledgerMock, brokerService); + + Method fence = PersistentTopic.class.getDeclaredMethod("fence"); + fence.setAccessible(true); + Method unfence = PersistentTopic.class.getDeclaredMethod("unfence"); + unfence.setAccessible(true); + + Field fencedTopicMonitoringTaskField = PersistentTopic.class.getDeclaredField("fencedTopicMonitoringTask"); + fencedTopicMonitoringTaskField.setAccessible(true); + Field isFencedField = AbstractTopic.class.getDeclaredField("isFenced"); + isFencedField.setAccessible(true); + Field isClosingOrDeletingField = PersistentTopic.class.getDeclaredField("isClosingOrDeleting"); + isClosingOrDeletingField.setAccessible(true); + + doReturn(10).when(svcConfig).getTopicFencingTimeoutSeconds(); + fence.invoke(topic); + unfence.invoke(topic); + ScheduledFuture fencedTopicMonitoringTask = (ScheduledFuture) fencedTopicMonitoringTaskField.get(topic); + assertTrue(fencedTopicMonitoringTask.isDone()); + assertTrue(fencedTopicMonitoringTask.isCancelled()); + assertFalse((boolean) isFencedField.get(topic)); + assertFalse((boolean) isClosingOrDeletingField.get(topic)); + + doReturn(1).when(svcConfig).getTopicFencingTimeoutSeconds(); + fence.invoke(topic); + Thread.sleep(2000); + fencedTopicMonitoringTask = (ScheduledFuture) fencedTopicMonitoringTaskField.get(topic); + assertTrue(fencedTopicMonitoringTask.isDone()); + assertFalse(fencedTopicMonitoringTask.isCancelled()); + assertTrue((boolean) isFencedField.get(topic)); + assertTrue((boolean) isClosingOrDeletingField.get(topic)); + } + private ByteBuf getMessageWithMetadata(byte[] data) throws IOException { MessageMetadata messageData = MessageMetadata.newBuilder().setPublishTime(System.currentTimeMillis()) .setProducerName("prod-name").setSequenceId(0).build();