From 64859d510051a0af07cf6ea2127478a52565a07a Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Mon, 18 Jul 2022 19:00:26 -0700 Subject: [PATCH 1/3] PulsarLedgerManager to pass correct error code to BK client; handle deletion of already deleted ledger --- .../impl/ManagedLedgerFactoryImpl.java | 27 +++- .../service/BrokerBkEnsemblesTests.java | 148 ++++++++++++++++++ .../metadata/api/MetadataStoreException.java | 46 ++++-- 3 files changed, 210 insertions(+), 11 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 1a12f9da49606..a851ce8131a7b 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -893,7 +893,18 @@ private void deleteManagedLedgerData(BookKeeper bkc, String managedLedgerName, M DeleteLedgerCallback callback, Object ctx) { Futures.waitForAll(info.ledgers.stream() .filter(li -> !li.isOffloaded) - .map(li -> bkc.newDeleteLedgerOp().withLedgerId(li.ledgerId).execute()) + .map(li -> bkc.newDeleteLedgerOp().withLedgerId(li.ledgerId).execute() + .handleAsync((result, ex) -> { + if (ex != null) { + int rc = BKException.getExceptionCode(ex); + if (rc == BKException.Code.NoSuchLedgerExistsOnMetadataServerException + || rc == BKException.Code.NoSuchLedgerExistsException) { + return null; + } + throw new CompletionException(ex); + } + return result; + })) .collect(Collectors.toList())) .thenRun(() -> { // Delete the metadata @@ -921,7 +932,19 @@ private CompletableFuture deleteCursor(BookKeeper bkc, String managedLedge // Delete the cursor ledger if present if (cursor.cursorsLedgerId != -1) { - cursorLedgerDeleteFuture = bkc.newDeleteLedgerOp().withLedgerId(cursor.cursorsLedgerId).execute(); + cursorLedgerDeleteFuture = bkc.newDeleteLedgerOp().withLedgerId(cursor.cursorsLedgerId) + .execute() + .handleAsync((result, ex) -> { + if (ex != null) { + int rc = BKException.getExceptionCode(ex); + if (rc == BKException.Code.NoSuchLedgerExistsOnMetadataServerException + || rc == BKException.Code.NoSuchLedgerExistsException) { + return null; + } + throw new CompletionException(ex); + } + return result; + }); } else { cursorLedgerDeleteFuture = CompletableFuture.completedFuture(null); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java index 8db3734eabeb1..fe215c4f00c3c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java @@ -20,6 +20,7 @@ import static org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest.retryStrategically; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.fail; import java.lang.reflect.Field; @@ -31,9 +32,13 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import com.google.common.collect.Sets; import lombok.Cleanup; +import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.client.BKException; import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.client.api.LedgerMetadata; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; @@ -47,12 +52,14 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.apache.zookeeper.ZooKeeper; import org.awaitility.Awaitility; import org.testng.Assert; import org.testng.annotations.Test; @Test(groups = "broker") +@Slf4j public class BrokerBkEnsemblesTests extends BkEnsemblesTestBase { public BrokerBkEnsemblesTests() { @@ -279,6 +286,147 @@ public void testSkipCorruptDataLedger() throws Exception { consumer.close(); } + @Test + public void testTruncateCorruptDataLedger() throws Exception { + // Ensure intended state for autoSkipNonRecoverableData + admin.brokers().updateDynamicConfiguration("autoSkipNonRecoverableData", "false"); + + @Cleanup + PulsarClient client = PulsarClient.builder() + .serviceUrl(pulsar.getWebServiceAddress()) + .statsInterval(0, TimeUnit.SECONDS) + .build(); + + final int totalMessages = 100; + final int totalDataLedgers = 5; + final int entriesPerLedger = totalMessages / totalDataLedgers; + + final String tenant = "prop"; + try { + admin.tenants().createTenant(tenant, new TenantInfoImpl(Sets.newHashSet("role1", "role2"), + Sets.newHashSet(config.getClusterName()))); + } catch (Exception e) { + + } + final String ns1 = tenant + "/crash-broker"; + try { + admin.namespaces().createNamespace(ns1, Sets.newHashSet(config.getClusterName())); + } catch (Exception e) { + + } + + final String topic1 = "persistent://" + ns1 + "/my-topic-" + System.currentTimeMillis(); + + // Create subscription + Consumer consumer = client.newConsumer().topic(topic1).subscriptionName("my-subscriber-name") + .receiverQueueSize(5).subscribe(); + + PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topic1).get(); + ManagedLedgerImpl ml = (ManagedLedgerImpl) topic.getManagedLedger(); + ManagedCursorImpl cursor = (ManagedCursorImpl) ml.getCursors().iterator().next(); + Field configField = ManagedCursorImpl.class.getDeclaredField("config"); + configField.setAccessible(true); + // Create multiple data-ledger + ManagedLedgerConfig config = (ManagedLedgerConfig) configField.get(cursor); + config.setMaxEntriesPerLedger(entriesPerLedger); + config.setMinimumRolloverTime(1, TimeUnit.MILLISECONDS); + // bookkeeper client + Field bookKeeperField = ManagedLedgerImpl.class.getDeclaredField("bookKeeper"); + bookKeeperField.setAccessible(true); + // Create multiple data-ledger + BookKeeper bookKeeper = (BookKeeper) bookKeeperField.get(ml); + + // (1) publish messages in 10 data-ledgers each with 20 entries under managed-ledger + Producer producer = client.newProducer().topic(topic1).create(); + for (int i = 0; i < totalMessages; i++) { + String message = "my-message-" + i; + producer.send(message.getBytes()); + } + + // validate: consumer is able to consume msg and close consumer after reading 1 entry + Assert.assertNotNull(consumer.receive(1, TimeUnit.SECONDS)); + consumer.close(); + + NavigableMap ledgerInfo = ml.getLedgersInfo(); + Assert.assertEquals(ledgerInfo.size(), totalDataLedgers); + Entry lastLedger = ledgerInfo.lastEntry(); + long firstLedgerToDelete = lastLedger.getKey(); + + // (2) delete first 4 data-ledgers + ledgerInfo.entrySet().forEach(entry -> { + if (!entry.equals(lastLedger)) { + assertEquals(entry.getValue().getEntries(), entriesPerLedger); + try { + bookKeeper.deleteLedger(entry.getKey()); + } catch (Exception e) { + e.printStackTrace(); + } + } + }); + + // create 5 more ledgers + for (int i = 0; i < totalMessages; i++) { + String message = "my-message2-" + i; + producer.send(message.getBytes()); + } + + ml.delete(); + + // Admin should be able to truncate the topic + admin.topics().truncate(topic1); + + ledgerInfo.entrySet().forEach(entry -> { + log.warn("found ledger: {}", entry.getKey()); + assertNotEquals(firstLedgerToDelete, entry.getKey()); + }); + + // Currently, ledger deletion is async and failed deletion + // does not actually fail truncation but logs an exception + // and creates scheduled task to retry + Awaitility.await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> { + LedgerMetadata meta = bookKeeper + .getLedgerMetadata(firstLedgerToDelete) + .exceptionally(e -> null) + .get(); + assertEquals(null, meta, "ledger should be deleted " + firstLedgerToDelete); + }); + + // Should not throw, deleting absent ledger must be a noop + // unless PulsarManager returned a wrong error which + // got translated to BKUnexpectedConditionException + try { + bookKeeper.deleteLedger(firstLedgerToDelete); + } catch (BKException.BKNoSuchLedgerExistsOnMetadataServerException bke) { + // pass + } + + producer.close(); + consumer.close(); + } + + @Test + public void testDeleteLedgerFactoryCorruptLedger() throws Exception { + ManagedLedgerFactoryImpl factory = (ManagedLedgerFactoryImpl) pulsar.getManagedLedgerFactory(); + ManagedLedgerImpl ml = (ManagedLedgerImpl) factory.open("test"); + + // bookkeeper client + Field bookKeeperField = ManagedLedgerImpl.class.getDeclaredField("bookKeeper"); + bookKeeperField.setAccessible(true); + // Create multiple data-ledger + BookKeeper bookKeeper = (BookKeeper) bookKeeperField.get(ml); + + ml.addEntry("dummy-entry-1".getBytes()); + + NavigableMap ledgerInfo = ml.getLedgersInfo(); + long lastLedger = ledgerInfo.lastEntry().getKey(); + + ml.close(); + bookKeeper.deleteLedger(lastLedger); + + // BK ledger is deleted, factory should not throw on delete + factory.delete("test"); + } + @Test(timeOut = 20000) public void testTopicWithWildCardChar() throws Exception { @Cleanup diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java index f1e72c27a6172..b7e3cc1346601 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java @@ -21,12 +21,31 @@ import java.io.IOException; import java.util.concurrent.CompletionException; import java.util.concurrent.ExecutionException; +import org.apache.bookkeeper.client.BKException; /** * Generic metadata store exception. */ public class MetadataStoreException extends IOException { + private static Throwable makeBkFriendlyException(int code, Throwable cause) { + Throwable t = cause; + while (t != null) { + if (t instanceof BKException) { + // BKException.getExceptionCode() traverses chain of causes + // no need to create another exception + return cause; + } + t = t.getCause(); + } + + BKException bkException = BKException.create(code); + if (cause != null) { + bkException.initCause(cause); + } + return bkException; + } + public MetadataStoreException(Throwable t) { super(t); } @@ -61,15 +80,18 @@ public InvalidImplementationException(String msg) { */ public static class NotFoundException extends MetadataStoreException { public NotFoundException() { - super((Throwable) null); + super(makeBkFriendlyException( + BKException.Code.NoSuchLedgerExistsOnMetadataServerException, null)); } public NotFoundException(Throwable t) { - super(t); + super(makeBkFriendlyException( + BKException.Code.NoSuchLedgerExistsOnMetadataServerException, t)); } public NotFoundException(String msg) { - super(msg); + super(msg, makeBkFriendlyException( + BKException.Code.NoSuchLedgerExistsOnMetadataServerException, null)); } } @@ -78,11 +100,13 @@ public NotFoundException(String msg) { */ public static class AlreadyExistsException extends MetadataStoreException { public AlreadyExistsException(Throwable t) { - super(t); + super(makeBkFriendlyException( + BKException.Code.LedgerExistException, t)); } public AlreadyExistsException(String msg) { - super(msg); + super(msg, makeBkFriendlyException( + BKException.Code.LedgerExistException, null)); } } @@ -91,11 +115,13 @@ public AlreadyExistsException(String msg) { */ public static class BadVersionException extends MetadataStoreException { public BadVersionException(Throwable t) { - super(t); + super(makeBkFriendlyException( + BKException.Code.MetadataVersionException, t)); } public BadVersionException(String msg) { - super(msg); + super(msg, makeBkFriendlyException( + BKException.Code.MetadataVersionException, null)); } } @@ -138,11 +164,13 @@ public LockBusyException(String msg) { */ public static class AlreadyClosedException extends MetadataStoreException { public AlreadyClosedException(Throwable t) { - super(t); + super(makeBkFriendlyException( + BKException.Code.LedgerClosedException, t)); } public AlreadyClosedException(String msg) { - super(msg); + super(msg, makeBkFriendlyException( + BKException.Code.LedgerClosedException, null)); } } From 8850e45aa7d5e9a7449a962fa13a1731390cb366 Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Mon, 18 Jul 2022 21:51:25 -0700 Subject: [PATCH 2/3] To not break code that relies on specific chain of exceptions, like ex.getCause().getCause() instanceof Blah --- .../metadata/api/MetadataStoreException.java | 21 ++++++++++++------- 1 file changed, 13 insertions(+), 8 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java index b7e3cc1346601..7044d875ff575 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java @@ -29,21 +29,26 @@ public class MetadataStoreException extends IOException { private static Throwable makeBkFriendlyException(int code, Throwable cause) { - Throwable t = cause; - while (t != null) { - if (t instanceof BKException) { + if (cause == null) { + return BKException.create(code); + } + + Throwable lastCause = cause; + while (true) { + if (lastCause instanceof BKException) { // BKException.getExceptionCode() traverses chain of causes // no need to create another exception return cause; } - t = t.getCause(); + if (lastCause.getCause() == null) { + break; + } + lastCause = lastCause.getCause(); } BKException bkException = BKException.create(code); - if (cause != null) { - bkException.initCause(cause); - } - return bkException; + lastCause.initCause(bkException); + return cause; } public MetadataStoreException(Throwable t) { From 40df0c4e5e9c76ed88a595afb861757dcffd25d0 Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Tue, 19 Jul 2022 15:50:16 -0700 Subject: [PATCH 3/3] CR feedback --- .../impl/ManagedLedgerFactoryImpl.java | 6 ++- .../service/BrokerBkEnsemblesTests.java | 4 +- .../metadata/api/MetadataStoreException.java | 41 +++++++++---------- 3 files changed, 26 insertions(+), 25 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index a851ce8131a7b..4bb3d18cffddc 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -894,11 +894,12 @@ private void deleteManagedLedgerData(BookKeeper bkc, String managedLedgerName, M Futures.waitForAll(info.ledgers.stream() .filter(li -> !li.isOffloaded) .map(li -> bkc.newDeleteLedgerOp().withLedgerId(li.ledgerId).execute() - .handleAsync((result, ex) -> { + .handle((result, ex) -> { if (ex != null) { int rc = BKException.getExceptionCode(ex); if (rc == BKException.Code.NoSuchLedgerExistsOnMetadataServerException || rc == BKException.Code.NoSuchLedgerExistsException) { + log.info("Ledger {} does not exist, ignoring", li.ledgerId); return null; } throw new CompletionException(ex); @@ -934,11 +935,12 @@ private CompletableFuture deleteCursor(BookKeeper bkc, String managedLedge if (cursor.cursorsLedgerId != -1) { cursorLedgerDeleteFuture = bkc.newDeleteLedgerOp().withLedgerId(cursor.cursorsLedgerId) .execute() - .handleAsync((result, ex) -> { + .handle((result, ex) -> { if (ex != null) { int rc = BKException.getExceptionCode(ex); if (rc == BKException.Code.NoSuchLedgerExistsOnMetadataServerException || rc == BKException.Code.NoSuchLedgerExistsException) { + log.info("Ledger {} does not exist, ignoring", cursor.cursorsLedgerId); return null; } throw new CompletionException(ex); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java index fe215c4f00c3c..612d9368b8c70 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java @@ -245,7 +245,7 @@ public void testSkipCorruptDataLedger() throws Exception { try { bookKeeper.deleteLedger(entry.getKey()); } catch (Exception e) { - e.printStackTrace(); + log.warn("failed to delete ledger {}", entry.getKey(), e); } } }); @@ -359,7 +359,7 @@ public void testTruncateCorruptDataLedger() throws Exception { try { bookKeeper.deleteLedger(entry.getKey()); } catch (Exception e) { - e.printStackTrace(); + log.warn("failed to delete ledger {}", entry.getKey(), e); } } }); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java index 7044d875ff575..05360f6a7d5ef 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/api/MetadataStoreException.java @@ -28,9 +28,9 @@ */ public class MetadataStoreException extends IOException { - private static Throwable makeBkFriendlyException(int code, Throwable cause) { + private static Throwable makeBkFriendlyException(BKException bkException, Throwable cause) { if (cause == null) { - return BKException.create(code); + return bkException; } Throwable lastCause = cause; @@ -46,7 +46,6 @@ private static Throwable makeBkFriendlyException(int code, Throwable cause) { lastCause = lastCause.getCause(); } - BKException bkException = BKException.create(code); lastCause.initCause(bkException); return cause; } @@ -84,19 +83,19 @@ public InvalidImplementationException(String msg) { * Key not found in store. */ public static class NotFoundException extends MetadataStoreException { + private static final BKException bkEx = + BKException.create(BKException.Code.NoSuchLedgerExistsOnMetadataServerException); + public NotFoundException() { - super(makeBkFriendlyException( - BKException.Code.NoSuchLedgerExistsOnMetadataServerException, null)); + super(makeBkFriendlyException(bkEx, null)); } public NotFoundException(Throwable t) { - super(makeBkFriendlyException( - BKException.Code.NoSuchLedgerExistsOnMetadataServerException, t)); + super(makeBkFriendlyException(bkEx, t)); } public NotFoundException(String msg) { - super(msg, makeBkFriendlyException( - BKException.Code.NoSuchLedgerExistsOnMetadataServerException, null)); + super(msg, makeBkFriendlyException(bkEx, null)); } } @@ -104,14 +103,14 @@ public NotFoundException(String msg) { * Key was already in store. */ public static class AlreadyExistsException extends MetadataStoreException { + private static final BKException bkEx = BKException.create(BKException.Code.LedgerExistException); + public AlreadyExistsException(Throwable t) { - super(makeBkFriendlyException( - BKException.Code.LedgerExistException, t)); + super(makeBkFriendlyException(bkEx, t)); } public AlreadyExistsException(String msg) { - super(msg, makeBkFriendlyException( - BKException.Code.LedgerExistException, null)); + super(msg, makeBkFriendlyException(bkEx, null)); } } @@ -119,14 +118,14 @@ public AlreadyExistsException(String msg) { * Unsuccessful update due to mismatched expected version. */ public static class BadVersionException extends MetadataStoreException { + private static final BKException bkEx = BKException.create(BKException.Code.MetadataVersionException); + public BadVersionException(Throwable t) { - super(makeBkFriendlyException( - BKException.Code.MetadataVersionException, t)); + super(makeBkFriendlyException(bkEx, t)); } public BadVersionException(String msg) { - super(msg, makeBkFriendlyException( - BKException.Code.MetadataVersionException, null)); + super(msg, makeBkFriendlyException(bkEx, null)); } } @@ -168,14 +167,14 @@ public LockBusyException(String msg) { * The store was already closed. */ public static class AlreadyClosedException extends MetadataStoreException { + private static final BKException bkEx = BKException.create(BKException.Code.LedgerClosedException); + public AlreadyClosedException(Throwable t) { - super(makeBkFriendlyException( - BKException.Code.LedgerClosedException, t)); + super(makeBkFriendlyException(bkEx, t)); } public AlreadyClosedException(String msg) { - super(msg, makeBkFriendlyException( - BKException.Code.LedgerClosedException, null)); + super(msg, makeBkFriendlyException(bkEx, null)); } }