From f8a75b4bd622e673f75464a8863eaf81d87209aa Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Thu, 14 Oct 2021 21:22:56 +0800 Subject: [PATCH 1/5] store and checkout --- .../service/persistent/PersistentTopic.java | 4 +++ .../broker/transaction/TransactionTest.java | 28 +++++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 0f403a87ca130..3b0804004e476 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -2987,6 +2987,10 @@ public void publishTxnMessage(TxnID txnID, ByteBuf headersAndPayload, PublishCon decrementPendingWriteOpsAndCheck(); }) .exceptionally(throwable -> { + throwable = throwable.getCause(); + if (!(throwable instanceof ManagedLedgerException)) { + throwable = new ManagedLedgerException(throwable); + } addFailed((ManagedLedgerException) throwable, publishContext); return null; }); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index fb94638bf49c1..398a1e0f91b2d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -18,15 +18,20 @@ */ package org.apache.pulsar.broker.transaction; +import static java.nio.charset.StandardCharsets.UTF_8; import static org.apache.pulsar.transaction.coordinator.impl.MLTransactionLogImpl.TRANSACTION_LOG_PREFIX; import com.google.common.collect.Sets; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; import java.lang.reflect.Field; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; @@ -34,6 +39,7 @@ import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.broker.transaction.buffer.impl.TopicTransactionBuffer; import org.apache.pulsar.broker.transaction.pendingack.PendingAckStore; import org.apache.pulsar.broker.transaction.pendingack.impl.MLPendingAckStore; import org.apache.pulsar.broker.transaction.pendingack.impl.MLPendingAckStoreProvider; @@ -246,4 +252,26 @@ public void testSubscriptionRecreateTopic() } + @Test + public void testAppendBufferWithManageLedgerException() + throws PulsarAdminException, ExecutionException, InterruptedException { + + String topic = "persistent://pulsar/system/testReCreateTopic"; + admin.topics().createNonPartitionedTopic(topic); + + PersistentTopic persistentTopic = + (PersistentTopic) pulsarServiceList.get(0).getBrokerService() + .getTopic(topic, false) + .get().get(); + + TopicTransactionBuffer topicTransactionBuffer = (TopicTransactionBuffer) persistentTopic.getTransactionBuffer(); + try { + topicTransactionBuffer.appendBufferToTxn(new TxnID(123L, 321L), 123L, + null); + Assert.fail("TransactionBuffer should not append successfully with a null buffer."); + } catch (Exception e) { + Assert.assertEquals(e.getClass(), ManagedLedgerException.class); + } + + } } \ No newline at end of file From 7aedf139cb74a4da9fba95ff51f8561bbd1cc937 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Fri, 15 Oct 2021 08:57:46 +0800 Subject: [PATCH 2/5] Add test --- .../broker/transaction/TransactionTest.java | 78 ++++++++++++++++--- 1 file changed, 67 insertions(+), 11 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 398a1e0f91b2d..0308d5cc378cc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -27,8 +27,10 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedLedgerException; @@ -252,10 +254,12 @@ public void testSubscriptionRecreateTopic() } + // If TXB appends buffer failed, it maybe return a exception that is not ManageLedgerException. + //But `future`.exceptionally need a ManageLedgerException as parameter for addFailed() @Test - public void testAppendBufferWithManageLedgerException() - throws PulsarAdminException, ExecutionException, InterruptedException { - + public void testAppendBufferWithNotManageLedgerExceptionCanCastToMLE() + throws PulsarAdminException, ExecutionException, InterruptedException, NoSuchFieldException, + IllegalAccessException { String topic = "persistent://pulsar/system/testReCreateTopic"; admin.topics().createNonPartitionedTopic(topic); @@ -263,15 +267,67 @@ public void testAppendBufferWithManageLedgerException() (PersistentTopic) pulsarServiceList.get(0).getBrokerService() .getTopic(topic, false) .get().get(); + CountDownLatch countDownLatch = new CountDownLatch(1); + //Make transactionBufferFuture complete. + Topic.PublishContext publishContext = new Topic.PublishContext() { - TopicTransactionBuffer topicTransactionBuffer = (TopicTransactionBuffer) persistentTopic.getTransactionBuffer(); - try { - topicTransactionBuffer.appendBufferToTxn(new TxnID(123L, 321L), 123L, - null); - Assert.fail("TransactionBuffer should not append successfully with a null buffer."); - } catch (Exception e) { - Assert.assertEquals(e.getClass(), ManagedLedgerException.class); - } + @Override + public String getProducerName() { + return "test"; + } + + public long getSequenceId() { + return 30; + } + /** + * Return the producer name for the original producer. + * + * For messages published locally, this will return the same local producer name, though in case of replicated + * messages, the original producer name will differ + */ + public String getOriginalProducerName() { + return "test"; + } + public long getOriginalSequenceId() { + return 30; + } + + public long getHighestSequenceId() { + return 30; + } + + public long getOriginalHighestSequenceId() { + return 30; + } + + public long getNumberOfMessages() { + return 30; + } + + @Override + public void completed(Exception e, long ledgerId, long entryId) { + Assert.assertTrue(e.getCause() instanceof ManagedLedgerException); + countDownLatch.countDown(); + } + }; + persistentTopic.publishTxnMessage(new TxnID(123L, 321L), + Unpooled.copiedBuffer("message", UTF_8), publishContext); + //Close manageLedger. + Class managedLedgerClass = ManagedLedgerImpl.class; + Field STATE_UPDATER_Class = managedLedgerClass.getDeclaredField("STATE_UPDATER"); + STATE_UPDATER_Class.setAccessible(true); + AtomicReferenceFieldUpdater STATE_UPDATER = + (AtomicReferenceFieldUpdater) + STATE_UPDATER_Class.get(persistentTopic.getManagedLedger()); + STATE_UPDATER.set((ManagedLedgerImpl) persistentTopic.getManagedLedger(), ManagedLedgerImpl.State.Closed); + //Publish to a close managerLedger to test ManagerLedgerException. + persistentTopic.publishTxnMessage(new TxnID(123L, 321L), + Unpooled.copiedBuffer("message", UTF_8), publishContext); + //If timeout, it means the assertTrue in publishContext.completed is failed. + Awaitility.await().until(() -> { + countDownLatch.await(); + return true; + }); } } \ No newline at end of file From d5b064581734ed7f7de161843eddf389f731cb77 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Fri, 15 Oct 2021 20:37:29 +0800 Subject: [PATCH 3/5] Optimization notes --- .../apache/pulsar/broker/transaction/TransactionTest.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 91e86388e16bf..f52debd671050 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -300,8 +300,7 @@ public void testTakeSnapshotBeforeBuildTxnProducer() throws Exception { }); } - //If TB appends buffer failed, it maybe return a exception that is not a ManageLedgerException. - //But `future`.exceptionally need a ManageLedgerException as parameter for addFailed() + @Test public void testAppendBufferWithNotManageLedgerExceptionCanCastToMLE() throws Exception { @@ -359,11 +358,11 @@ public void completed(Exception e, long ledgerId, long entryId) { //Close topic manageLedger. persistentTopic.getManagedLedger().close(); - //Publish to a close managerLedger to test ManagerLedgerException. + //Publish to a closed managerLedger to test ManagerLedgerException. persistentTopic.publishTxnMessage(new TxnID(123L, 321L), Unpooled.copiedBuffer("message", UTF_8), publishContext); - //If timeout, it means the assertTrue in publishContext.completed is failed. + //If it times out, it means that the assertTrue in publishContext.completed is failed. Awaitility.await().until(() -> { countDownLatch.await(); return true; From ef90bf10177b800e36a387de2d22561541dcc8f3 Mon Sep 17 00:00:00 2001 From: congbo Date: Sat, 16 Oct 2021 00:28:27 +0800 Subject: [PATCH 4/5] change publish transaction message return exception --- .../transaction/buffer/impl/TopicTransactionBuffer.java | 2 +- .../org/apache/pulsar/broker/transaction/TransactionTest.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) 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 79f7f35732754..1b0fb6697ec3f 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 @@ -219,7 +219,7 @@ public void addComplete(Position position, ByteBuf entryData, Object ctx) { @Override public void addFailed(ManagedLedgerException exception, Object ctx) { log.error("Failed to append buffer to txn {}", txnId, exception); - completableFuture.completeExceptionally(new PersistenceException(exception)); + completableFuture.completeExceptionally(exception); } }, null); return completableFuture; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index f52debd671050..d88efdeb31346 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -302,7 +302,7 @@ public void testTakeSnapshotBeforeBuildTxnProducer() throws Exception { @Test - public void testAppendBufferWithNotManageLedgerExceptionCanCastToMLE() + public void testPublishTransactionMessageThrowException() throws Exception { String topic = "persistent://pulsar/system/testReCreateTopic"; admin.topics().createNonPartitionedTopic(topic); @@ -350,7 +350,7 @@ public long getNumberOfMessages() { @Override public void completed(Exception e, long ledgerId, long entryId) { - Assert.assertTrue(e.getCause() instanceof ManagedLedgerException); + Assert.assertTrue(e.getCause() instanceof ManagedLedgerException.ManagedLedgerAlreadyClosedException); countDownLatch.countDown(); } }; From ff5ddb1dd98a9c00123f91507e23b58c139f7430 Mon Sep 17 00:00:00 2001 From: liangyepianzhou Date: Tue, 19 Oct 2021 16:37:27 +0800 Subject: [PATCH 5/5] merge --- .../org/apache/pulsar/broker/transaction/TransactionTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index d88efdeb31346..bb172b599c730 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -302,7 +302,7 @@ public void testTakeSnapshotBeforeBuildTxnProducer() throws Exception { @Test - public void testPublishTransactionMessageThrowException() + public void testAppendBufferWithNotManageLedgerExceptionCanCastToMLE() throws Exception { String topic = "persistent://pulsar/system/testReCreateTopic"; admin.topics().createNonPartitionedTopic(topic);