From fc32b1f1c1b046e4d1bf0a2d9473e160821eac57 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 6 Dec 2022 16:18:47 +0800 Subject: [PATCH 1/8] [improve] [test] Add test testTrimLedgerWillKeepsAtLeastOneLedgerWithData --- .../TopicTransactionBufferRecoverTest.java | 188 +++++++++++++++++- 1 file changed, 180 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java index 39c324d92f38c..6542571a86851 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java @@ -33,15 +33,19 @@ import java.lang.reflect.Field; import java.util.LinkedList; import java.util.List; +import java.util.Map; import java.util.NavigableMap; import java.util.Optional; +import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import lombok.AllArgsConstructor; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedLedgerException; +import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.impl.ReadOnlyManagedLedgerImpl; @@ -74,6 +78,7 @@ import org.apache.pulsar.client.api.Reader; import org.apache.pulsar.client.api.ReaderBuilder; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.client.impl.MessageIdImpl; @@ -119,18 +124,18 @@ protected void cleanup() throws Exception { } @DataProvider(name = "testTopic") - public Object[] testTopic() { - return new Object[] { - RECOVER_ABORT, - RECOVER_COMMIT + public Object[][] testTopic() { + return new Object[][] { + {RECOVER_ABORT}, + {RECOVER_COMMIT} }; } @DataProvider(name = "enableSnapshotSegment") - public Object[] testSnapshot() { - return new Boolean[] { - true, - false + public Object[][] testSnapshot() { + return new Object[][] { + {true}, + {false} }; } @@ -248,6 +253,173 @@ private void recoverTest(String testTopic) throws Exception { } + private ProducerAndConsumer makeManyTx(int txCount, String topicName, String subName) throws Exception { + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .subscriptionType(SubscriptionType.Shared) + .topic(topicName) + .isAckReceiptEnabled(true) + .acknowledgmentGroupTime(0, TimeUnit.SECONDS) + .subscriptionName(subName) + .subscribe(); + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .sendTimeout(0, TimeUnit.SECONDS) + .enableBatching(false) + .batchingMaxMessages(2) + .create(); + producer.send("first message"); + boolean lastTxCommitted = false; + Message lastMessage = null; + for(int i = 0; i < txCount; i++) { + Transaction transaction = + pulsarClient.newTransaction().withTransactionTimeout(10, TimeUnit.SECONDS).build().get(); + lastMessage = consumer.receive(); + producer.newMessage(transaction) + .value(new StringBuilder("tx message 0-") + .append(String.valueOf(lastMessage.getMessageId())).toString()).sendAsync(); + producer.newMessage(transaction) + .value(new StringBuilder("tx message 1-") + .append(String.valueOf(lastMessage.getMessageId())).toString()).sendAsync(); + consumer.acknowledgeAsync(lastMessage.getMessageId(), transaction); + if (i % 2 == 0) { + transaction.commit().get(); + lastTxCommitted = true; + } else { + transaction.abort().get(); + lastTxCommitted = false; + } + } + if (lastTxCommitted){ + Message msg = consumer.receive(); + consumer.acknowledge(msg); + } else { + consumer.acknowledge(lastMessage); + } + return new ProducerAndConsumer(producer, consumer); + } + + @AllArgsConstructor + private static class ProducerAndConsumer { + public Producer producer; + public Consumer consumer; + } + + private PersistentTopic findPersistentTopic(String topicName){ + for (PulsarService pulsarService : pulsarServiceList){ + CompletableFuture> future = pulsarService.getBrokerService().getTopic(topicName, false); + if (future == null || !future.isDone() || future.isCompletedExceptionally() || !future.join().isPresent()){ + continue; + } + return (PersistentTopic) future.join().get(); + } + throw new RuntimeException("topic[" + topicName + "] not found."); + } + + private void triggerSnapshot(String topicName){ + PersistentTopic persistentTopic = findPersistentTopic(topicName); + TopicTransactionBuffer topicTransactionBuffer = + (TopicTransactionBuffer) persistentTopic.getTransactionBuffer(); + topicTransactionBuffer.run(null); + } + + private void triggerLedgerTrims(String topicName){ + PersistentTopic persistentTopic = findPersistentTopic(topicName); + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + CompletableFuture future = new CompletableFuture(); + managedLedger.trimConsumedLedgersInBackground(future); + future.join(); + } + + private Map getLedgers(String topicName){ + PersistentTopic persistentTopic = findPersistentTopic(topicName); + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + return managedLedger.getLedgersInfo(); + } + + private void triggerCompact(String topicName) throws Exception { + PersistentTopic persistentTopic = findPersistentTopic(topicName); + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + persistentTopic.getBrokerService().getPulsar().getCompactor().compact(topicName); + Awaitility.await().untilAsserted(() -> { + ManagedCursorImpl compaction = (ManagedCursorImpl) managedLedger.getCursors().get("__compaction"); + assertEquals(compaction.getMarkDeletedPosition().getLedgerId(), + managedLedger.getLastConfirmedEntry().getLedgerId()); + assertEquals(compaction.getMarkDeletedPosition().getEntryId(), + managedLedger.getLastConfirmedEntry().getEntryId()); + }); + ManagedCursorImpl compaction = (ManagedCursorImpl) managedLedger.getCursors().get("__compaction"); + log.info("===> cursor-compaction mark deleted position {}:{}", compaction.getMarkDeletedPosition().getLedgerId(), + compaction.getMarkDeletedPosition().getEntryId()); + } + + private void waitCursorDedup(String topicName) throws Exception { + PersistentTopic persistentTopic = findPersistentTopic(topicName); + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + persistentTopic.checkDeduplicationSnapshot(); + Awaitility.await().untilAsserted(() -> { + ManagedCursorImpl dedupCursor = (ManagedCursorImpl) managedLedger.getCursors().get("pulsar.dedup"); + assertEquals(dedupCursor.getMarkDeletedPosition().getLedgerId(), + managedLedger.getLastConfirmedEntry().getLedgerId()); + assertEquals(dedupCursor.getMarkDeletedPosition().getEntryId(), + managedLedger.getLastConfirmedEntry().getEntryId()); + }); + ManagedCursorImpl dedupCursor = (ManagedCursorImpl) managedLedger.getCursors().get("pulsar.dedup"); + log.info("===> cursor-dedup mark deleted position {}:{}", dedupCursor.getMarkDeletedPosition().getLedgerId(), + dedupCursor.getMarkDeletedPosition().getEntryId()); + } + + @Test + private void testTrimLedgerWillKeepsAtLeastOneLedgerWithData() throws Exception { + String topicName = String.format("persistent://%s/%s", NAMESPACE1, + "tx_recover_" + UUID.randomUUID().toString().replaceAll("-", "_")); + String subName = "sub"; + String transactionBufferTopicName = + String.format("persistent://%s/%s", NAMESPACE1, TRANSACTION_BUFFER_SNAPSHOT); + + // Make some data. + ProducerAndConsumer producerAndConsumer = null; + for (int i = 0; i < 5; i++) { + producerAndConsumer = makeManyTx(10, topicName, subName); + triggerSnapshot(topicName); + if (i != 4) { + // Do not close all clients. + producerAndConsumer.producer.close(); + producerAndConsumer.consumer.close(); + } + // Reload for create new ledger, and wait for topic reload. + admin.topics().unload(transactionBufferTopicName); + Awaitility.await().until(() -> { + try { + findPersistentTopic(transactionBufferTopicName); + return true; + } catch (Exception e) { + return false; + } + }); + } + + // Verify the last ledger will not be deleted. + Map ledgers = getLedgers(transactionBufferTopicName); + long lastLedgerHasData = -1; + for (MLDataFormats.ManagedLedgerInfo.LedgerInfo ledger : ledgers.values()){ + if (ledger.getEntries() > 0){ + lastLedgerHasData = Math.max(lastLedgerHasData, ledger.getLedgerId()); + } + } + log.info("===> ledgers before trim {}", ledgers.keySet()); + triggerCompact(transactionBufferTopicName); + waitCursorDedup(transactionBufferTopicName); + triggerLedgerTrims(transactionBufferTopicName); + ledgers = getLedgers(transactionBufferTopicName); + log.info("===> ledgers after trim {}", ledgers.keySet()); + assertTrue(ledgers.containsKey(lastLedgerHasData)); + + // cleanup. + producerAndConsumer.producer.close(); + producerAndConsumer.consumer.close(); + admin.topics().delete(topicName, false); + } + private void testTakeSnapshot() throws Exception { @Cleanup Producer producer = pulsarClient From ac8c5480355f447771852e0de4ad83fe7319aaf5 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 9 Dec 2022 16:27:43 +0800 Subject: [PATCH 2/8] [fix] [tx] Transaction buffer recover can not completed --- .../broker/systopic/SystemTopicClient.java | 11 ++ .../TopicPoliciesSystemTopicClient.java | 6 + ...onBufferSnapshotBaseSystemTopicClient.java | 6 + ...SingleSnapshotAbortedTxnProcessorImpl.java | 16 ++- .../buffer/impl/TopicTransactionBuffer.java | 7 +- .../impl/TopicTransactionBufferState.java | 4 + .../buffer/TransactionBufferException.java | 12 ++ .../TopicTransactionBufferRecoverTest.java | 128 ++++++++---------- 8 files changed, 116 insertions(+), 74 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java index 88ca099b4ca17..b71509075feaf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java @@ -21,6 +21,7 @@ import java.io.IOException; import java.util.List; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClientException; @@ -155,6 +156,16 @@ interface Reader { */ Message readNext() throws PulsarClientException; + /** + * Read the next message in the system topic waiting for a maximum time. + * + *

Returns null if no message is received before the timeout. + * + * @return the next message(Could be null if none received in time) + * @throws PulsarClientException + */ + Message readNext(int timeout, TimeUnit unit) throws PulsarClientException; + /** * Async read event from system topic. * @return pulsar event future diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java index 3fd8921c15efa..a8cb5f73d4536 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java @@ -22,6 +22,7 @@ import java.util.ArrayList; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; @@ -171,6 +172,11 @@ public Message readNext() throws PulsarClientException { return reader.readNext(); } + @Override + public Message readNext(int timeout, TimeUnit unit) throws PulsarClientException { + return reader.readNext(timeout, unit); + } + @Override public CompletableFuture> readNextAsync() { return reader.readNextAsync(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java index b18bf552c3004..801c40726c72b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java @@ -20,6 +20,7 @@ import java.io.IOException; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.service.SystemTopicTxnBufferSnapshotService; import org.apache.pulsar.client.api.Message; @@ -143,6 +144,11 @@ public Message readNext() throws PulsarClientException { return reader.readNext(); } + @Override + public Message readNext(int timeout, TimeUnit unit) throws PulsarClientException { + return reader.readNext(timeout, unit); + } + @Override public CompletableFuture> readNextAsync() { return reader.readNextAsync(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java index a13dd0499a6bc..30ccfd8a3dd2a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; @@ -31,8 +32,10 @@ import org.apache.pulsar.broker.transaction.buffer.AbortedTxnProcessor; import org.apache.pulsar.broker.transaction.buffer.metadata.AbortTxnMetadata; import org.apache.pulsar.broker.transaction.buffer.metadata.TransactionBufferSnapshot; +import org.apache.pulsar.broker.transaction.exception.buffer.TransactionBufferException; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.transaction.TxnID; +import org.apache.pulsar.common.naming.SystemTopicNames; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; @@ -87,7 +90,7 @@ public CompletableFuture recoverFromSnapshot() { PositionImpl startReadCursorPosition = null; try { while (reader.hasMoreEvents()) { - Message message = reader.readNext(); + Message message = reader.readNext(2, TimeUnit.SECONDS); if (topic.getName().equals(message.getKey())) { TransactionBufferSnapshot transactionBufferSnapshot = message.getValue(); if (transactionBufferSnapshot != null) { @@ -98,13 +101,20 @@ public CompletableFuture recoverFromSnapshot() { } } } - closeReader(reader); return CompletableFuture.completedFuture(startReadCursorPosition); + } catch (NullPointerException npe) { + String warn = String.format("[%s] When reading from topic %s,the latest message has been" + + " deleted by compaction-task or trim ledger.", + topic.getName(), + SystemTopicNames.TRANSACTION_BUFFER_SNAPSHOT); + log.warn(warn); + return FutureUtil.failedFuture(new TransactionBufferException.TBRecoverCantCompletedException(warn)); } catch (Exception ex) { log.error("[{}] Transaction buffer recover fail when read " + "transactionBufferSnapshot!", topic.getName(), ex); - closeReader(reader); return FutureUtil.failedFuture(ex); + } finally { + closeReader(reader); } }, topic.getBrokerService().getPulsar().getTransactionExecutorProvider() diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java index f3bf4f95923cd..1fde123b93eb5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java @@ -48,6 +48,7 @@ import org.apache.pulsar.broker.transaction.buffer.TransactionBufferReader; import org.apache.pulsar.broker.transaction.buffer.TransactionMeta; import org.apache.pulsar.broker.transaction.buffer.metadata.TransactionBufferSnapshot; +import org.apache.pulsar.broker.transaction.exception.buffer.TransactionBufferException; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.MessageMetadata; @@ -174,7 +175,11 @@ public void handleTxnEntry(Entry entry) { @Override public void recoverExceptionally(Throwable e) { - + if (e instanceof TransactionBufferException.TBRecoverCantCompletedException){ + changeToNoneState(); + recover(); + return; + } log.warn("Closing topic {} due to read transaction buffer snapshot while recovering the " + "transaction buffer throw exception", topic.getName(), e); // when create reader or writer fail throw PulsarClientException, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java index 92ab1d07b690d..f33e3932ee88d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java @@ -66,6 +66,10 @@ protected void changeToCloseState() { STATE_UPDATER.set(this, State.Close); } + protected void changeToNoneState() { + STATE_UPDATER.set(this, State.None); + } + public boolean checkIfReady() { return STATE_UPDATER.get(this) == State.Ready; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java index b1c4fdd1dbc25..99ccc5e62d862 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java @@ -87,5 +87,17 @@ public TransactionNotFoundException(String message) { } } + /** + * Exception is thrown when the transaction is not found in the transaction buffer. + */ + public static class TBRecoverCantCompletedException extends TransactionBufferException { + + private static final long serialVersionUID = 0L; + + public TBRecoverCantCompletedException(String message) { + super(message); + } + } + } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java index 6542571a86851..08721d401d7da 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java @@ -20,9 +20,11 @@ import static org.apache.pulsar.common.naming.SystemTopicNames.TRANSACTION_BUFFER_SNAPSHOT; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; @@ -33,12 +35,12 @@ import java.lang.reflect.Field; import java.util.LinkedList; import java.util.List; -import java.util.Map; import java.util.NavigableMap; import java.util.Optional; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import lombok.AllArgsConstructor; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; @@ -81,7 +83,9 @@ import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.client.api.transaction.TxnID; +import org.apache.pulsar.client.impl.ConsumerImpl; import org.apache.pulsar.client.impl.MessageIdImpl; +import org.apache.pulsar.client.impl.ReaderImpl; import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.events.EventType; @@ -91,6 +95,7 @@ import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; import org.awaitility.Awaitility; +import org.awaitility.reflect.WhiteboxImpl; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.DataProvider; @@ -322,21 +327,7 @@ private void triggerSnapshot(String topicName){ topicTransactionBuffer.run(null); } - private void triggerLedgerTrims(String topicName){ - PersistentTopic persistentTopic = findPersistentTopic(topicName); - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); - CompletableFuture future = new CompletableFuture(); - managedLedger.trimConsumedLedgersInBackground(future); - future.join(); - } - - private Map getLedgers(String topicName){ - PersistentTopic persistentTopic = findPersistentTopic(topicName); - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); - return managedLedger.getLedgersInfo(); - } - - private void triggerCompact(String topicName) throws Exception { + private void triggerCompactAndWait(String topicName) throws Exception { PersistentTopic persistentTopic = findPersistentTopic(topicName); ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); persistentTopic.getBrokerService().getPulsar().getCompactor().compact(topicName); @@ -352,71 +343,68 @@ private void triggerCompact(String topicName) throws Exception { compaction.getMarkDeletedPosition().getEntryId()); } - private void waitCursorDedup(String topicName) throws Exception { - PersistentTopic persistentTopic = findPersistentTopic(topicName); - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); - persistentTopic.checkDeduplicationSnapshot(); - Awaitility.await().untilAsserted(() -> { - ManagedCursorImpl dedupCursor = (ManagedCursorImpl) managedLedger.getCursors().get("pulsar.dedup"); - assertEquals(dedupCursor.getMarkDeletedPosition().getLedgerId(), - managedLedger.getLastConfirmedEntry().getLedgerId()); - assertEquals(dedupCursor.getMarkDeletedPosition().getEntryId(), - managedLedger.getLastConfirmedEntry().getEntryId()); - }); - ManagedCursorImpl dedupCursor = (ManagedCursorImpl) managedLedger.getCursors().get("pulsar.dedup"); - log.info("===> cursor-dedup mark deleted position {}:{}", dedupCursor.getMarkDeletedPosition().getLedgerId(), - dedupCursor.getMarkDeletedPosition().getEntryId()); + private void initPropLastMessageIdInBrokerOfTBReader(TopicName topicName) throws Exception { + for (PulsarService pulsarService : pulsarServiceList){ + // Init prop: lastMessageIdInBroker. + final SystemTopicTxnBufferSnapshotService tbSnapshotService = + pulsarService.getTransactionBufferSnapshotServiceFactory().getTxnBufferSnapshotService(); + AtomicReference readerCache = new AtomicReference<>(); + SystemTopicClient.Reader reader = + (SystemTopicClient.Reader) tbSnapshotService.createReader(topicName).join(); + ReaderImpl innerReader = + WhiteboxImpl.getInternalState(reader, "reader"); + ConsumerImpl innerConsumer = innerReader.getConsumer(); + // Because the consumer's incoming queue is 1000, messages are placed in client memory the moment consumers + // are created, which prevents the consumer from pulling messages from the broker again. + // So we just set a fake MessageId to variable "lastMessageIdInBroker". + innerConsumer.hasMessageAvailable(); + Field lastMessageInBrokerField = ConsumerImpl.class.getDeclaredField("lastMessageIdInBroker"); + lastMessageInBrokerField.setAccessible(true); + lastMessageInBrokerField.set(innerConsumer, MessageId.latest); + readerCache.set(reader); + // Inject. + SystemTopicTxnBufferSnapshotService spyTbSnapshotService = spy(tbSnapshotService); + doAnswer(invocation -> { + SystemTopicClient.Reader readerInCache = readerCache.getAndSet(null); + if (readerInCache != null){ + return CompletableFuture.completedFuture(readerInCache); + } else { + return tbSnapshotService.createReader(topicName); + } + }).when(spyTbSnapshotService).createReader(topicName); + Field field = + TransactionBufferSnapshotServiceFactory.class.getDeclaredField("txnBufferSnapshotService"); + field.setAccessible(true); + field.set(pulsarService.getTransactionBufferSnapshotServiceFactory(), spyTbSnapshotService); + } } @Test - private void testTrimLedgerWillKeepsAtLeastOneLedgerWithData() throws Exception { + public void testRecoverCantComplete() throws Exception { + String topicNameForMakeDataForTB = String.format("persistent://%s/%s", NAMESPACE1, + "tx_recover_1_" + UUID.randomUUID().toString().replaceAll("-", "_")); String topicName = String.format("persistent://%s/%s", NAMESPACE1, - "tx_recover_" + UUID.randomUUID().toString().replaceAll("-", "_")); + "tx_recover_2_" + UUID.randomUUID().toString().replaceAll("-", "_")); String subName = "sub"; String transactionBufferTopicName = String.format("persistent://%s/%s", NAMESPACE1, TRANSACTION_BUFFER_SNAPSHOT); - // Make some data. - ProducerAndConsumer producerAndConsumer = null; - for (int i = 0; i < 5; i++) { - producerAndConsumer = makeManyTx(10, topicName, subName); - triggerSnapshot(topicName); - if (i != 4) { - // Do not close all clients. - producerAndConsumer.producer.close(); - producerAndConsumer.consumer.close(); - } - // Reload for create new ledger, and wait for topic reload. - admin.topics().unload(transactionBufferTopicName); - Awaitility.await().until(() -> { - try { - findPersistentTopic(transactionBufferTopicName); - return true; - } catch (Exception e) { - return false; - } - }); - } + // Make some TB snapshot. + ProducerAndConsumer producerAndConsumer1 = makeManyTx(10, topicNameForMakeDataForTB, subName); + triggerSnapshot(topicNameForMakeDataForTB); + producerAndConsumer1.producer.close(); + producerAndConsumer1.consumer.close(); + admin.topics().delete(topicNameForMakeDataForTB); - // Verify the last ledger will not be deleted. - Map ledgers = getLedgers(transactionBufferTopicName); - long lastLedgerHasData = -1; - for (MLDataFormats.ManagedLedgerInfo.LedgerInfo ledger : ledgers.values()){ - if (ledger.getEntries() > 0){ - lastLedgerHasData = Math.max(lastLedgerHasData, ledger.getLedgerId()); - } - } - log.info("===> ledgers before trim {}", ledgers.keySet()); - triggerCompact(transactionBufferTopicName); - waitCursorDedup(transactionBufferTopicName); - triggerLedgerTrims(transactionBufferTopicName); - ledgers = getLedgers(transactionBufferTopicName); - log.info("===> ledgers after trim {}", ledgers.keySet()); - assertTrue(ledgers.containsKey(lastLedgerHasData)); + // Make race condition of "getLastMessageId" and "compaction" to make recover can't complete. + initPropLastMessageIdInBrokerOfTBReader(TopicName.get(topicName)); + triggerCompactAndWait(transactionBufferTopicName); + // Ensure topic works. + ProducerAndConsumer producerAndConsumer2 = makeManyTx(3, topicName, subName); // cleanup. - producerAndConsumer.producer.close(); - producerAndConsumer.consumer.close(); + producerAndConsumer2.producer.close(); + producerAndConsumer2.consumer.close(); admin.topics().delete(topicName, false); } From 0dd281805dbd8f60b02bb4bf4ba1fee6d43834f1 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 9 Dec 2022 17:12:39 +0800 Subject: [PATCH 3/8] hanle npe more precision --- .../SingleSnapshotAbortedTxnProcessorImpl.java | 17 +++++++++-------- 1 file changed, 9 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java index 30ccfd8a3dd2a..e077a04954b35 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java @@ -91,6 +91,15 @@ public CompletableFuture recoverFromSnapshot() { try { while (reader.hasMoreEvents()) { Message message = reader.readNext(2, TimeUnit.SECONDS); + if (message == null){ + String warnLog = String.format("[%s] When reading from topic %s,the latest message has" + + " been deleted by compaction-task or trim ledger.", + topic.getName(), + SystemTopicNames.TRANSACTION_BUFFER_SNAPSHOT); + log.warn(warnLog); + return FutureUtil.failedFuture( + new TransactionBufferException.TBRecoverCantCompletedException(warnLog)); + } if (topic.getName().equals(message.getKey())) { TransactionBufferSnapshot transactionBufferSnapshot = message.getValue(); if (transactionBufferSnapshot != null) { @@ -102,13 +111,6 @@ public CompletableFuture recoverFromSnapshot() { } } return CompletableFuture.completedFuture(startReadCursorPosition); - } catch (NullPointerException npe) { - String warn = String.format("[%s] When reading from topic %s,the latest message has been" - + " deleted by compaction-task or trim ledger.", - topic.getName(), - SystemTopicNames.TRANSACTION_BUFFER_SNAPSHOT); - log.warn(warn); - return FutureUtil.failedFuture(new TransactionBufferException.TBRecoverCantCompletedException(warn)); } catch (Exception ex) { log.error("[{}] Transaction buffer recover fail when read " + "transactionBufferSnapshot!", topic.getName(), ex); @@ -116,7 +118,6 @@ public CompletableFuture recoverFromSnapshot() { } finally { closeReader(reader); } - }, topic.getBrokerService().getPulsar().getTransactionExecutorProvider() .getExecutor(this)); } From cdff5e180c361ab0f1a02689f621c654819cc2db Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Sun, 11 Dec 2022 02:35:38 +0800 Subject: [PATCH 4/8] avoid unnecessary changes --- .../TopicTransactionBufferRecoverTest.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java index 08721d401d7da..17a8baf98116e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java @@ -129,18 +129,18 @@ protected void cleanup() throws Exception { } @DataProvider(name = "testTopic") - public Object[][] testTopic() { - return new Object[][] { - {RECOVER_ABORT}, - {RECOVER_COMMIT} + public Object[] testTopic() { + return new Object[] { + RECOVER_ABORT, + RECOVER_COMMIT }; } @DataProvider(name = "enableSnapshotSegment") - public Object[][] testSnapshot() { - return new Object[][] { - {true}, - {false} + public Object[] testSnapshot() { + return new Boolean[] { + true, + false }; } From 73f90ef46253024c482471015d185511d6132b9d Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Mon, 12 Dec 2022 14:51:55 +0800 Subject: [PATCH 5/8] Simplified logic --- .../broker/systopic/SystemTopicClient.java | 11 ---------- .../TopicPoliciesSystemTopicClient.java | 6 ------ ...onBufferSnapshotBaseSystemTopicClient.java | 6 ------ ...SingleSnapshotAbortedTxnProcessorImpl.java | 20 +++++++++---------- .../buffer/impl/TopicTransactionBuffer.java | 7 +------ .../impl/TopicTransactionBufferState.java | 4 ---- .../buffer/TransactionBufferException.java | 12 ----------- 7 files changed, 10 insertions(+), 56 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java index b71509075feaf..88ca099b4ca17 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/SystemTopicClient.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.util.List; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.PulsarClientException; @@ -156,16 +155,6 @@ interface Reader { */ Message readNext() throws PulsarClientException; - /** - * Read the next message in the system topic waiting for a maximum time. - * - *

Returns null if no message is received before the timeout. - * - * @return the next message(Could be null if none received in time) - * @throws PulsarClientException - */ - Message readNext(int timeout, TimeUnit unit) throws PulsarClientException; - /** * Async read event from system topic. * @return pulsar event future diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java index a8cb5f73d4536..3fd8921c15efa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TopicPoliciesSystemTopicClient.java @@ -22,7 +22,6 @@ import java.util.ArrayList; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; @@ -172,11 +171,6 @@ public Message readNext() throws PulsarClientException { return reader.readNext(); } - @Override - public Message readNext(int timeout, TimeUnit unit) throws PulsarClientException { - return reader.readNext(timeout, unit); - } - @Override public CompletableFuture> readNextAsync() { return reader.readNextAsync(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java index 801c40726c72b..b18bf552c3004 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/systopic/TransactionBufferSnapshotBaseSystemTopicClient.java @@ -20,7 +20,6 @@ import java.io.IOException; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.service.SystemTopicTxnBufferSnapshotService; import org.apache.pulsar.client.api.Message; @@ -144,11 +143,6 @@ public Message readNext() throws PulsarClientException { return reader.readNext(); } - @Override - public Message readNext(int timeout, TimeUnit unit) throws PulsarClientException { - return reader.readNext(timeout, unit); - } - @Override public CompletableFuture> readNextAsync() { return reader.readNextAsync(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java index e077a04954b35..f92c465c01672 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java @@ -32,10 +32,8 @@ import org.apache.pulsar.broker.transaction.buffer.AbortedTxnProcessor; import org.apache.pulsar.broker.transaction.buffer.metadata.AbortTxnMetadata; import org.apache.pulsar.broker.transaction.buffer.metadata.TransactionBufferSnapshot; -import org.apache.pulsar.broker.transaction.exception.buffer.TransactionBufferException; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.transaction.TxnID; -import org.apache.pulsar.common.naming.SystemTopicNames; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; @@ -90,15 +88,15 @@ public CompletableFuture recoverFromSnapshot() { PositionImpl startReadCursorPosition = null; try { while (reader.hasMoreEvents()) { - Message message = reader.readNext(2, TimeUnit.SECONDS); - if (message == null){ - String warnLog = String.format("[%s] When reading from topic %s,the latest message has" - + " been deleted by compaction-task or trim ledger.", - topic.getName(), - SystemTopicNames.TRANSACTION_BUFFER_SNAPSHOT); - log.warn(warnLog); - return FutureUtil.failedFuture( - new TransactionBufferException.TBRecoverCantCompletedException(warnLog)); + Message message; + try { + message = reader.readNextAsync().get(30, TimeUnit.SECONDS); + } catch (Exception ex) { + Throwable t = FutureUtil.unwrapCompletionException(ex); + log.error("[{}] Transaction buffer recover fail when read " + + "transactionBufferSnapshot!", topic.getName(), t); + closeReader(reader); + return FutureUtil.failedFuture(t); } if (topic.getName().equals(message.getKey())) { TransactionBufferSnapshot transactionBufferSnapshot = message.getValue(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java index 1fde123b93eb5..f3bf4f95923cd 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java @@ -48,7 +48,6 @@ import org.apache.pulsar.broker.transaction.buffer.TransactionBufferReader; import org.apache.pulsar.broker.transaction.buffer.TransactionMeta; import org.apache.pulsar.broker.transaction.buffer.metadata.TransactionBufferSnapshot; -import org.apache.pulsar.broker.transaction.exception.buffer.TransactionBufferException; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.common.api.proto.MessageMetadata; @@ -175,11 +174,7 @@ public void handleTxnEntry(Entry entry) { @Override public void recoverExceptionally(Throwable e) { - if (e instanceof TransactionBufferException.TBRecoverCantCompletedException){ - changeToNoneState(); - recover(); - return; - } + log.warn("Closing topic {} due to read transaction buffer snapshot while recovering the " + "transaction buffer throw exception", topic.getName(), e); // when create reader or writer fail throw PulsarClientException, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java index f33e3932ee88d..92ab1d07b690d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBufferState.java @@ -66,10 +66,6 @@ protected void changeToCloseState() { STATE_UPDATER.set(this, State.Close); } - protected void changeToNoneState() { - STATE_UPDATER.set(this, State.None); - } - public boolean checkIfReady() { return STATE_UPDATER.get(this) == State.Ready; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java index 99ccc5e62d862..b1c4fdd1dbc25 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/exception/buffer/TransactionBufferException.java @@ -87,17 +87,5 @@ public TransactionNotFoundException(String message) { } } - /** - * Exception is thrown when the transaction is not found in the transaction buffer. - */ - public static class TBRecoverCantCompletedException extends TransactionBufferException { - - private static final long serialVersionUID = 0L; - - public TBRecoverCantCompletedException(String message) { - super(message); - } - } - } From 171878205605d8d2c51d231005aa807e355de71d Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 13 Dec 2022 11:52:53 +0800 Subject: [PATCH 6/8] fix test --- .../buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java | 6 +++--- .../transaction/TopicTransactionBufferRecoverTest.java | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java index f92c465c01672..1cd001f2803c3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java @@ -22,6 +22,7 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; @@ -90,12 +91,11 @@ public CompletableFuture recoverFromSnapshot() { while (reader.hasMoreEvents()) { Message message; try { - message = reader.readNextAsync().get(30, TimeUnit.SECONDS); - } catch (Exception ex) { + message = reader.readNextAsync().get(5, TimeUnit.SECONDS); + } catch (TimeoutException ex) { Throwable t = FutureUtil.unwrapCompletionException(ex); log.error("[{}] Transaction buffer recover fail when read " + "transactionBufferSnapshot!", topic.getName(), t); - closeReader(reader); return FutureUtil.failedFuture(t); } if (topic.getName().equals(message.getKey())) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java index 17a8baf98116e..cc3eef9ba9e81 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java @@ -379,7 +379,7 @@ private void initPropLastMessageIdInBrokerOfTBReader(TopicName topicName) throws } } - @Test + @Test(timeOut = 60 * 1000) public void testRecoverCantComplete() throws Exception { String topicNameForMakeDataForTB = String.format("persistent://%s/%s", NAMESPACE1, "tx_recover_1_" + UUID.randomUUID().toString().replaceAll("-", "_")); From 51d1f64fc5f3eec236891e0aedf2985a7f0fc649 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 14 Dec 2022 15:55:52 +0800 Subject: [PATCH 7/8] address comments --- ...SingleSnapshotAbortedTxnProcessorImpl.java | 25 ++++++++++++------- .../TopicTransactionBufferRecoverTest.java | 6 ++--- 2 files changed, 18 insertions(+), 13 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java index 1cd001f2803c3..97c315573b54a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java @@ -28,6 +28,7 @@ import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.collections4.map.LinkedMap; +import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.systopic.SystemTopicClient; import org.apache.pulsar.broker.transaction.buffer.AbortedTxnProcessor; @@ -35,6 +36,7 @@ import org.apache.pulsar.broker.transaction.buffer.metadata.TransactionBufferSnapshot; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.transaction.TxnID; +import org.apache.pulsar.client.impl.PulsarClientImpl; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; @@ -80,24 +82,22 @@ public boolean checkAbortedTransaction(TxnID txnID, Position readPosition) { return aborts.containsKey(txnID); } + private long getSystemClientOperationTimeoutMs() throws Exception { + PulsarClientImpl pulsarClient = (PulsarClientImpl) topic.getBrokerService().getPulsar().getClient(); + return pulsarClient.getConfiguration().getOperationTimeoutMs(); + } @Override public CompletableFuture recoverFromSnapshot() { return topic.getBrokerService().getPulsar().getTransactionBufferSnapshotServiceFactory() .getTxnBufferSnapshotService() .createReader(TopicName.get(topic.getName())).thenComposeAsync(reader -> { - PositionImpl startReadCursorPosition = null; try { + PositionImpl startReadCursorPosition = null; while (reader.hasMoreEvents()) { Message message; - try { - message = reader.readNextAsync().get(5, TimeUnit.SECONDS); - } catch (TimeoutException ex) { - Throwable t = FutureUtil.unwrapCompletionException(ex); - log.error("[{}] Transaction buffer recover fail when read " - + "transactionBufferSnapshot!", topic.getName(), t); - return FutureUtil.failedFuture(t); - } + message = reader.readNextAsync().get(getSystemClientOperationTimeoutMs(), + TimeUnit.MILLISECONDS); if (topic.getName().equals(message.getKey())) { TransactionBufferSnapshot transactionBufferSnapshot = message.getValue(); if (transactionBufferSnapshot != null) { @@ -109,6 +109,13 @@ public CompletableFuture recoverFromSnapshot() { } } return CompletableFuture.completedFuture(startReadCursorPosition); + } catch (TimeoutException ex) { + Throwable t = FutureUtil.unwrapCompletionException(ex); + String errorMessage = String.format("[%s] Transaction buffer recover fail by read " + + "transactionBufferSnapshot timeout!", topic.getName()); + log.error(errorMessage, t); + return FutureUtil.failedFuture( + new BrokerServiceException.ServiceUnitNotReadyException(errorMessage, t)); } catch (Exception ex) { log.error("[{}] Transaction buffer recover fail when read " + "transactionBufferSnapshot!", topic.getName(), ex); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java index cc3eef9ba9e81..2c4d7df7ec856 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java @@ -114,6 +114,7 @@ public class TopicTransactionBufferRecoverTest extends TransactionTestBase { private static final int NUM_PARTITIONS = 16; @BeforeMethod protected void setup() throws Exception { + conf.getProperties().setProperty("brokerClient_operationTimeoutMs", Integer.valueOf(10 * 1000).toString()); setUpBase(1, NUM_PARTITIONS, RECOVER_COMMIT, 0); admin.topics().createNonPartitionedTopic(RECOVER_ABORT); admin.topics().createNonPartitionedTopic(TAKE_SNAPSHOT); @@ -333,10 +334,7 @@ private void triggerCompactAndWait(String topicName) throws Exception { persistentTopic.getBrokerService().getPulsar().getCompactor().compact(topicName); Awaitility.await().untilAsserted(() -> { ManagedCursorImpl compaction = (ManagedCursorImpl) managedLedger.getCursors().get("__compaction"); - assertEquals(compaction.getMarkDeletedPosition().getLedgerId(), - managedLedger.getLastConfirmedEntry().getLedgerId()); - assertEquals(compaction.getMarkDeletedPosition().getEntryId(), - managedLedger.getLastConfirmedEntry().getEntryId()); + assertEquals(compaction.getMarkDeletedPosition(), managedLedger.getLastConfirmedEntry()); }); ManagedCursorImpl compaction = (ManagedCursorImpl) managedLedger.getCursors().get("__compaction"); log.info("===> cursor-compaction mark deleted position {}:{}", compaction.getMarkDeletedPosition().getLedgerId(), From 01d5cec7dff3b0128a1342dd469ebbcddb4bbe96 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Dec 2022 15:14:48 +0800 Subject: [PATCH 8/8] make test simpler --- ...SingleSnapshotAbortedTxnProcessorImpl.java | 5 +- .../TopicTransactionBufferRecoverTest.java | 177 +++++------------- 2 files changed, 52 insertions(+), 130 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java index 97c315573b54a..f8d0d32391233 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/SingleSnapshotAbortedTxnProcessorImpl.java @@ -95,9 +95,8 @@ public CompletableFuture recoverFromSnapshot() { try { PositionImpl startReadCursorPosition = null; while (reader.hasMoreEvents()) { - Message message; - message = reader.readNextAsync().get(getSystemClientOperationTimeoutMs(), - TimeUnit.MILLISECONDS); + Message message = reader.readNextAsync() + .get(getSystemClientOperationTimeoutMs(), TimeUnit.MILLISECONDS); if (topic.getName().equals(message.getKey())) { TransactionBufferSnapshot transactionBufferSnapshot = message.getValue(); if (transactionBufferSnapshot != null) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java index 2c4d7df7ec856..a2b72fc458db4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TopicTransactionBufferRecoverTest.java @@ -25,6 +25,7 @@ import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; @@ -40,14 +41,12 @@ import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicReference; -import lombok.AllArgsConstructor; +import java.util.concurrent.atomic.AtomicBoolean; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedLedgerException; -import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.bookkeeper.mledger.impl.ReadOnlyManagedLedgerImpl; @@ -80,12 +79,9 @@ import org.apache.pulsar.client.api.Reader; import org.apache.pulsar.client.api.ReaderBuilder; import org.apache.pulsar.client.api.Schema; -import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.client.api.transaction.TxnID; -import org.apache.pulsar.client.impl.ConsumerImpl; import org.apache.pulsar.client.impl.MessageIdImpl; -import org.apache.pulsar.client.impl.ReaderImpl; import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.events.EventType; @@ -95,7 +91,6 @@ import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; import org.awaitility.Awaitility; -import org.awaitility.reflect.WhiteboxImpl; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.DataProvider; @@ -259,117 +254,54 @@ private void recoverTest(String testTopic) throws Exception { } - private ProducerAndConsumer makeManyTx(int txCount, String topicName, String subName) throws Exception { - Consumer consumer = pulsarClient.newConsumer(Schema.STRING) - .subscriptionType(SubscriptionType.Shared) - .topic(topicName) - .isAckReceiptEnabled(true) - .acknowledgmentGroupTime(0, TimeUnit.SECONDS) - .subscriptionName(subName) - .subscribe(); - Producer producer = pulsarClient.newProducer(Schema.STRING) - .topic(topicName) - .sendTimeout(0, TimeUnit.SECONDS) - .enableBatching(false) - .batchingMaxMessages(2) - .create(); - producer.send("first message"); - boolean lastTxCommitted = false; - Message lastMessage = null; - for(int i = 0; i < txCount; i++) { - Transaction transaction = - pulsarClient.newTransaction().withTransactionTimeout(10, TimeUnit.SECONDS).build().get(); - lastMessage = consumer.receive(); - producer.newMessage(transaction) - .value(new StringBuilder("tx message 0-") - .append(String.valueOf(lastMessage.getMessageId())).toString()).sendAsync(); - producer.newMessage(transaction) - .value(new StringBuilder("tx message 1-") - .append(String.valueOf(lastMessage.getMessageId())).toString()).sendAsync(); - consumer.acknowledgeAsync(lastMessage.getMessageId(), transaction); - if (i % 2 == 0) { - transaction.commit().get(); - lastTxCommitted = true; + private void makeTBSnapshotReaderTimeoutIfFirstRead(TopicName topicName) throws Exception { + SystemTopicClient.Reader mockReader = mock(SystemTopicClient.Reader.class); + AtomicBoolean isFirstCallOfMethodHasMoreEvents = new AtomicBoolean(); + AtomicBoolean isFirstCallOfMethodHasReadNext = new AtomicBoolean(); + AtomicBoolean isFirstCallOfMethodHasReadNextAsync = new AtomicBoolean(); + + doAnswer(invocation -> { + if (isFirstCallOfMethodHasMoreEvents.compareAndSet(false,true)){ + return true; } else { - transaction.abort().get(); - lastTxCommitted = false; + return false; } - } - if (lastTxCommitted){ - Message msg = consumer.receive(); - consumer.acknowledge(msg); - } else { - consumer.acknowledge(lastMessage); - } - return new ProducerAndConsumer(producer, consumer); - } + }).when(mockReader).hasMoreEvents(); - @AllArgsConstructor - private static class ProducerAndConsumer { - public Producer producer; - public Consumer consumer; - } - - private PersistentTopic findPersistentTopic(String topicName){ - for (PulsarService pulsarService : pulsarServiceList){ - CompletableFuture> future = pulsarService.getBrokerService().getTopic(topicName, false); - if (future == null || !future.isDone() || future.isCompletedExceptionally() || !future.join().isPresent()){ - continue; + doAnswer(invocation -> { + if (isFirstCallOfMethodHasReadNext.compareAndSet(false, true)){ + // Just stuck the thread. + Thread.sleep(3600 * 1000); } - return (PersistentTopic) future.join().get(); - } - throw new RuntimeException("topic[" + topicName + "] not found."); - } - - private void triggerSnapshot(String topicName){ - PersistentTopic persistentTopic = findPersistentTopic(topicName); - TopicTransactionBuffer topicTransactionBuffer = - (TopicTransactionBuffer) persistentTopic.getTransactionBuffer(); - topicTransactionBuffer.run(null); - } + return null; + }).when(mockReader).readNext(); + + doAnswer(invocation -> { + CompletableFuture future = new CompletableFuture<>(); + new Thread(() -> { + if (isFirstCallOfMethodHasReadNextAsync.compareAndSet(false, true)){ + // Just stuck the thread. + try { + Thread.sleep(3600 * 1000); + } catch (InterruptedException e) { + } + future.complete(null); + } else { + future.complete(null); + } + }).start(); + return future; + }).when(mockReader).readNextAsync(); - private void triggerCompactAndWait(String topicName) throws Exception { - PersistentTopic persistentTopic = findPersistentTopic(topicName); - ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); - persistentTopic.getBrokerService().getPulsar().getCompactor().compact(topicName); - Awaitility.await().untilAsserted(() -> { - ManagedCursorImpl compaction = (ManagedCursorImpl) managedLedger.getCursors().get("__compaction"); - assertEquals(compaction.getMarkDeletedPosition(), managedLedger.getLastConfirmedEntry()); - }); - ManagedCursorImpl compaction = (ManagedCursorImpl) managedLedger.getCursors().get("__compaction"); - log.info("===> cursor-compaction mark deleted position {}:{}", compaction.getMarkDeletedPosition().getLedgerId(), - compaction.getMarkDeletedPosition().getEntryId()); - } + when(mockReader.closeAsync()).thenReturn(CompletableFuture.completedFuture(null)); - private void initPropLastMessageIdInBrokerOfTBReader(TopicName topicName) throws Exception { for (PulsarService pulsarService : pulsarServiceList){ // Init prop: lastMessageIdInBroker. final SystemTopicTxnBufferSnapshotService tbSnapshotService = pulsarService.getTransactionBufferSnapshotServiceFactory().getTxnBufferSnapshotService(); - AtomicReference readerCache = new AtomicReference<>(); - SystemTopicClient.Reader reader = - (SystemTopicClient.Reader) tbSnapshotService.createReader(topicName).join(); - ReaderImpl innerReader = - WhiteboxImpl.getInternalState(reader, "reader"); - ConsumerImpl innerConsumer = innerReader.getConsumer(); - // Because the consumer's incoming queue is 1000, messages are placed in client memory the moment consumers - // are created, which prevents the consumer from pulling messages from the broker again. - // So we just set a fake MessageId to variable "lastMessageIdInBroker". - innerConsumer.hasMessageAvailable(); - Field lastMessageInBrokerField = ConsumerImpl.class.getDeclaredField("lastMessageIdInBroker"); - lastMessageInBrokerField.setAccessible(true); - lastMessageInBrokerField.set(innerConsumer, MessageId.latest); - readerCache.set(reader); - // Inject. SystemTopicTxnBufferSnapshotService spyTbSnapshotService = spy(tbSnapshotService); - doAnswer(invocation -> { - SystemTopicClient.Reader readerInCache = readerCache.getAndSet(null); - if (readerInCache != null){ - return CompletableFuture.completedFuture(readerInCache); - } else { - return tbSnapshotService.createReader(topicName); - } - }).when(spyTbSnapshotService).createReader(topicName); + doAnswer(invocation -> CompletableFuture.completedFuture(mockReader)) + .when(spyTbSnapshotService).createReader(topicName); Field field = TransactionBufferSnapshotServiceFactory.class.getDeclaredField("txnBufferSnapshotService"); field.setAccessible(true); @@ -378,31 +310,22 @@ private void initPropLastMessageIdInBrokerOfTBReader(TopicName topicName) throws } @Test(timeOut = 60 * 1000) - public void testRecoverCantComplete() throws Exception { - String topicNameForMakeDataForTB = String.format("persistent://%s/%s", NAMESPACE1, - "tx_recover_1_" + UUID.randomUUID().toString().replaceAll("-", "_")); + public void testTBRecoverCanRetryIfTimeoutRead() throws Exception { String topicName = String.format("persistent://%s/%s", NAMESPACE1, - "tx_recover_2_" + UUID.randomUUID().toString().replaceAll("-", "_")); - String subName = "sub"; - String transactionBufferTopicName = - String.format("persistent://%s/%s", NAMESPACE1, TRANSACTION_BUFFER_SNAPSHOT); - - // Make some TB snapshot. - ProducerAndConsumer producerAndConsumer1 = makeManyTx(10, topicNameForMakeDataForTB, subName); - triggerSnapshot(topicNameForMakeDataForTB); - producerAndConsumer1.producer.close(); - producerAndConsumer1.consumer.close(); - admin.topics().delete(topicNameForMakeDataForTB); + "tx_recover_" + UUID.randomUUID().toString().replaceAll("-", "_")); // Make race condition of "getLastMessageId" and "compaction" to make recover can't complete. - initPropLastMessageIdInBrokerOfTBReader(TopicName.get(topicName)); - triggerCompactAndWait(transactionBufferTopicName); - // Ensure topic works. - ProducerAndConsumer producerAndConsumer2 = makeManyTx(3, topicName, subName); + makeTBSnapshotReaderTimeoutIfFirstRead(TopicName.get(topicName)); + // Verify( Cmd-PRODUCER will wait for TB recover finished ) + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topicName) + .sendTimeout(0, TimeUnit.SECONDS) + .enableBatching(false) + .batchingMaxMessages(2) + .create(); // cleanup. - producerAndConsumer2.producer.close(); - producerAndConsumer2.consumer.close(); + producer.close(); admin.topics().delete(topicName, false); }