From ee2fc2594b064d4cc263c54513c440d547e9c4e0 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Fri, 29 May 2020 20:08:35 +0800 Subject: [PATCH 1/3] add feature: enable rollover when triggering maxLedgerRolloverTimeMinutes --- .../bookkeeper/mledger/ManagedLedger.java | 5 + .../mledger/impl/ManagedLedgerImpl.java | 33 ++++++ .../pulsar/broker/service/BrokerService.java | 22 ++++ .../CurrentLedgerRolloverIfFullTest.java | 100 ++++++++++++++++++ 4 files changed, 160 insertions(+) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 714dd0eb6403f..eacacb56bda23 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -472,4 +472,9 @@ void asyncSetProperties(Map properties, final AsyncCallbacks.Set * @param promise */ void trimConsumedLedgersInBackground(CompletableFuture promise); + + /** + * Roll current ledger if it is full + */ + void rollCurrentLedgerIfFull(); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index e70b60ca9635f..fb99c84d78526 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -68,6 +68,7 @@ import java.util.function.Supplier; import java.util.stream.Collectors; +import org.apache.bookkeeper.client.AsyncCallback; import org.apache.bookkeeper.client.AsyncCallback.CreateCallback; import org.apache.bookkeeper.client.AsyncCallback.OpenCallback; import org.apache.bookkeeper.client.BKException; @@ -1391,6 +1392,38 @@ synchronized void ledgerClosed(final LedgerHandle lh) { } } + synchronized void createLedgerAfterClosed() { + STATE_UPDATER.set(this, State.CreatingLedger); + this.lastLedgerCreationInitiationTimestamp = System.nanoTime(); + mbean.startDataLedgerCreateOp(); + asyncCreateLedger(bookKeeper, config, digestType, this, Collections.emptyMap()); + } + + @Override + public void rollCurrentLedgerIfFull() { + log.info("[{}] Start checking if current ledger is full", name); + if (currentLedgerEntries > 0 && currentLedgerIsFull()) { + STATE_UPDATER.set(this, State.ClosingLedger); + currentLedger.asyncClose(new AsyncCallback.CloseCallback() { + @Override + public void closeComplete(int rc, LedgerHandle lh, Object o) { + checkArgument(currentLedger.getId() == lh.getId(), "ledgerId %s doesn't match with acked ledgerId %s", + currentLedger.getId(), + lh.getId()); + + if (rc == BKException.Code.OK) { + log.debug("Successfuly closed ledger {}", lh.getId()); + } else { + log.warn("Error when closing ledger {}. Status={}", lh.getId(), BKException.getMessage(rc)); + } + + ledgerClosed(lh); + createLedgerAfterClosed(); + } + }, System.nanoTime()); + } + } + void clearPendingAddEntries(ManagedLedgerException e) { while (!pendingAddEntries.isEmpty()) { OpAddEntry op = pendingAddEntries.poll(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 70015235f465b..1aa2d75f22118 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -194,6 +194,7 @@ public class BrokerService implements Closeable, ZooKeeperCacheListener { + if (t instanceof PersistentTopic) { + Optional.ofNullable(((PersistentTopic) t).getManagedLedger()).ifPresent( + managedLedger -> { + managedLedger.rollCurrentLedgerIfFull(); + } + ); + } + }); + } + public void checkMessageDeduplicationInfo() { forEachTopic(Topic::checkMessageDeduplicationInfo); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java new file mode 100644 index 0000000000000..8f102bff1640a --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java @@ -0,0 +1,100 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.service; + +import lombok.Cleanup; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; +import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.Producer; +import org.junit.Test; +import org.testng.Assert; + +import java.util.concurrent.TimeUnit; + +public class CurrentLedgerRolloverIfFullTest extends BrokerTestBase { + @Override + protected void setup() throws Exception { + + } + + @Override + protected void cleanup() throws Exception { + + } + + @Test + public void testCurrentLedgerRolloverIfFull() throws Exception { + conf.setRetentionCheckIntervalInSeconds(1); + conf.setManagedLedgerMaxLedgerRolloverTimeMinutes(1); + super.baseSetup(); + final String topicName = "persistent://prop/ns-abc/CurrentLedgerRolloverIfFullTest"; + final String subscriptionName = "CurrentLedgerRolloverIfFullTest-subscriber-name"; + + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(topicName) + .producerName("CurrentLedgerRolloverIfFullTest-producer-name") + .create(); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName) + .subscribe(); + + Topic topicRef = pulsar.getBrokerService().getTopicReference(topicName).get(); + Assert.assertNotNull(topicRef); + PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); + + ManagedLedgerConfig managedLedgerConfig = persistentTopic.getManagedLedger().getConfig(); + managedLedgerConfig.setRetentionSizeInMB(1L); + managedLedgerConfig.setRetentionTime(1, TimeUnit.SECONDS); + managedLedgerConfig.setMaxEntriesPerLedger(1); + managedLedgerConfig.setMinimumRolloverTime(1, TimeUnit.MILLISECONDS); +// managedLedgerConfig.setMaximumRolloverTime(3000, TimeUnit.MILLISECONDS); + + int msgNum = 5; + for (int i = 0; i < msgNum; i++) { + producer.send(new byte[1024 * 1024]); + } + + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), msgNum); + + for (int i = 0; i < msgNum; i++) { + Message msg = consumer.receive(2, TimeUnit.SECONDS); + Assert.assertTrue(msg != null); + consumer.acknowledge(msg); + } + + // all the messages reached the TTL + // and all the ledgers will be removed except the the last ledger + Thread.sleep(2500); + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1); + Assert.assertNotEquals(managedLedger.getCurrentLedgerSize(), 0); + + // the last ledger reached the max time interval of rollover + // and all the entries will be removed + Thread.sleep(65000); + managedLedger.rollCurrentLedgerIfFull(); + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1); + Assert.assertEquals(managedLedger.getCurrentLedgerSize(), 0); + } +} From b1558940ab2a0ef515845fe78b5de5ba9ce38ad7 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Sat, 30 May 2020 11:44:23 +0800 Subject: [PATCH 2/3] add supplemental test cases --- .../CurrentLedgerRolloverIfFullTest.java | 35 ++++++++++--------- 1 file changed, 19 insertions(+), 16 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java index 8f102bff1640a..b70a595341c03 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java @@ -21,6 +21,7 @@ import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.bookkeeper.mledger.util.Futures; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; @@ -43,11 +44,8 @@ protected void cleanup() throws Exception { @Test public void testCurrentLedgerRolloverIfFull() throws Exception { - conf.setRetentionCheckIntervalInSeconds(1); - conf.setManagedLedgerMaxLedgerRolloverTimeMinutes(1); super.baseSetup(); final String topicName = "persistent://prop/ns-abc/CurrentLedgerRolloverIfFullTest"; - final String subscriptionName = "CurrentLedgerRolloverIfFullTest-subscriber-name"; @Cleanup Producer producer = pulsarClient.newProducer() @@ -56,7 +54,9 @@ public void testCurrentLedgerRolloverIfFull() throws Exception { .create(); @Cleanup - Consumer consumer = pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName) + Consumer consumer = pulsarClient.newConsumer() + .topic(topicName) + .subscriptionName("CurrentLedgerRolloverIfFullTest-subscriber-name") .subscribe(); Topic topicRef = pulsar.getBrokerService().getTopicReference(topicName).get(); @@ -64,19 +64,18 @@ public void testCurrentLedgerRolloverIfFull() throws Exception { PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); ManagedLedgerConfig managedLedgerConfig = persistentTopic.getManagedLedger().getConfig(); - managedLedgerConfig.setRetentionSizeInMB(1L); managedLedgerConfig.setRetentionTime(1, TimeUnit.SECONDS); - managedLedgerConfig.setMaxEntriesPerLedger(1); + managedLedgerConfig.setMaxEntriesPerLedger(2); managedLedgerConfig.setMinimumRolloverTime(1, TimeUnit.MILLISECONDS); -// managedLedgerConfig.setMaximumRolloverTime(3000, TimeUnit.MILLISECONDS); + managedLedgerConfig.setMaximumRolloverTime(5, TimeUnit.MILLISECONDS); - int msgNum = 5; + int msgNum = 10; for (int i = 0; i < msgNum; i++) { producer.send(new byte[1024 * 1024]); } ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); - Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), msgNum); + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), msgNum / 2); for (int i = 0; i < msgNum; i++) { Message msg = consumer.receive(2, TimeUnit.SECONDS); @@ -84,17 +83,21 @@ public void testCurrentLedgerRolloverIfFull() throws Exception { consumer.acknowledge(msg); } - // all the messages reached the TTL - // and all the ledgers will be removed except the the last ledger - Thread.sleep(2500); + // all the messages have been acknowledged + // and all the ledgers have been removed except the the last ledger + Thread.sleep(500); Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1); Assert.assertNotEquals(managedLedger.getCurrentLedgerSize(), 0); - // the last ledger reached the max time interval of rollover - // and all the entries will be removed - Thread.sleep(65000); + // trigger a ledger rollover + // and now we have two ledgers, one with expired data and one for empty managedLedger.rollCurrentLedgerIfFull(); - Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1); + Thread.sleep(1000); + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 2); + + // trigger a ledger trimming + // and now we only have the empty ledger + managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); Assert.assertEquals(managedLedger.getCurrentLedgerSize(), 0); } } From 33f0615f4e120d92531dc080f725dcc25203fbb7 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Mon, 1 Jun 2020 11:53:54 +0800 Subject: [PATCH 3/3] shutdown monitor when BrokerService is closed --- .../java/org/apache/pulsar/broker/service/BrokerService.java | 1 + 1 file changed, 1 insertion(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 1aa2d75f22118..5116fe6f9b06b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -627,6 +627,7 @@ public void close() throws IOException { inactivityMonitor.shutdown(); messageExpiryMonitor.shutdown(); compactionMonitor.shutdown(); + ledgerFullMonitor.shutdown(); backlogQuotaChecker.shutdown(); authenticationService.close(); pulsarStats.close();