From 2255cffc331f49cb13081b8ba82464aa0f0e20b4 Mon Sep 17 00:00:00 2001 From: Masahiro Sakamoto Date: Fri, 13 Nov 2020 17:28:29 +0900 Subject: [PATCH 1/3] Close topics that remain fenced forcefully --- conf/broker.conf | 4 ++ conf/standalone.conf | 4 ++ .../pulsar/broker/ServiceConfiguration.java | 6 +++ .../service/persistent/PersistentTopic.java | 41 ++++++++++++++++++- .../broker/service/PersistentTopicTest.java | 38 +++++++++++++++++ 5 files changed, 91 insertions(+), 2 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 16688069caa13..0bcd00ce1e6b9 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 0 or less, 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..957a822b16044 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 0 or less, 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..1931b995dca3a 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 0 or less, 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..cbde6035d6cf5 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 synchronized 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(); From 5c495685d1028c01bf79cf11ae90e1aef687533b Mon Sep 17 00:00:00 2001 From: Masahiro Sakamoto Date: Sat, 14 Nov 2020 00:24:17 +0900 Subject: [PATCH 2/3] Fix comments --- conf/broker.conf | 2 +- conf/standalone.conf | 2 +- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 0bcd00ce1e6b9..5fe9e8dccdb9c 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -453,7 +453,7 @@ systemTopicEnabled=false topicLevelPoliciesEnabled=false # If a topic remains fenced for this number of seconds, it will be closed forcefully. -# If it is 0 or less, the fenced topic will not be closed. +# If it is set to 0 or a negative number, the fenced topic will not be closed. topicFencingTimeoutSeconds=0 ### --- Authentication --- ### diff --git a/conf/standalone.conf b/conf/standalone.conf index 957a822b16044..4cea72d82c128 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -366,7 +366,7 @@ systemTopicEnabled=false topicLevelPoliciesEnabled=false # If a topic remains fenced for this number of seconds, it will be closed forcefully. -# If it is 0 or less, the fenced topic will not be closed. +# If it is set to 0 or a negative number, the fenced topic will not be closed. topicFencingTimeoutSeconds=0 ### --- Authentication --- ### 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 1931b995dca3a..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 @@ -830,7 +830,7 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext( category = CATEGORY_SERVER, doc = "If a topic remains fenced for this number of seconds, it will be closed forcefully.\n" - + " If it is 0 or less, the fenced topic will not be closed." + + " If it is set to 0 or a negative number, the fenced topic will not be closed." ) private int topicFencingTimeoutSeconds = 0; From 2032c994af7ff9d84256c0605a85d42b2ad7c706 Mon Sep 17 00:00:00 2001 From: Masahiro Sakamoto Date: Tue, 17 Nov 2020 11:33:12 +0900 Subject: [PATCH 3/3] Remove synchronized from closeFencedTopicForcefully method --- .../pulsar/broker/service/persistent/PersistentTopic.java | 2 +- 1 file changed, 1 insertion(+), 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 cbde6035d6cf5..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 @@ -2407,7 +2407,7 @@ private synchronized void unfence() { } } - private synchronized void closeFencedTopicForcefully() { + private void closeFencedTopicForcefully() { if (isFenced) { final int timeout = brokerService.pulsar().getConfiguration().getTopicFencingTimeoutSeconds(); if (isClosingOrDeleting) {