From dfe8e7088cd4cbfd84c9f53d30698d68dfbe8c6e Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Thu, 12 Jan 2023 12:14:26 +0800 Subject: [PATCH 1/9] fix timeout transaction. --- .../apache/pulsar/broker/TransactionMetadataStoreService.java | 1 + 1 file changed, 1 insertion(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index 3d9e6924d1168..aca266fcbd202 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -448,6 +448,7 @@ private CompletableFuture endTxnInTransactionBuffer(TxnID txnID, int txnAc private static boolean isRetryableException(Throwable ex) { Throwable realCause = FutureUtil.unwrapCompletionException(ex); return (realCause instanceof TransactionMetadataStoreStateException + || realCause instanceof CoordinatorNotFoundException || realCause instanceof RequestTimeoutException || realCause instanceof ManagedLedgerException || realCause instanceof BrokerPersistenceException From 3ff9ceb72ee46e0cf557c39edd34b8425bda3d12 Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 16 Jan 2023 12:24:22 +0800 Subject: [PATCH 2/9] change patch way. --- .../TransactionMetadataStoreService.java | 81 ++++++++++--------- 1 file changed, 41 insertions(+), 40 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index aca266fcbd202..a88998b7325cd 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -133,54 +133,56 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc return; } - openTransactionMetadataStore(tcId).thenAccept((store) -> internalPinnedExecutor.execute(() -> { + openTransactionMetadataStore(tcId).thenAccept((store) -> { stores.put(tcId, store); LOG.info("Added new transaction meta store {}", tcId); - long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; - while (true) { - // prevent thread in a busy loop. - if (System.currentTimeMillis() < endTime) { - CompletableFuture future = deque.poll(); - if (future != null) { - // complete queue request future - future.complete(null); + internalPinnedExecutor.execute(() -> { + long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; + while (true) { + // prevent thread in a busy loop. + if (System.currentTimeMillis() < endTime) { + CompletableFuture future = deque.poll(); + if (future != null) { + // complete queue request future + future.complete(null); + } else { + break; + } } else { + deque.clear(); break; } - } else { - deque.clear(); - break; } - } - completableFuture.complete(null); - tcLoadSemaphore.release(); - })).exceptionally(e -> { + completableFuture.complete(null); + tcLoadSemaphore.release(); + }); + }).exceptionally(e -> { internalPinnedExecutor.execute(() -> { - completableFuture.completeExceptionally(e.getCause()); - // release before handle request queue, - //in order to client reconnect infinite loop - tcLoadSemaphore.release(); - long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; - while (true) { - // prevent thread in a busy loop. - if (System.currentTimeMillis() < endTime) { - CompletableFuture future = deque.poll(); - if (future != null) { - // this means that this tc client connection connect fail - future.completeExceptionally(e); - } else { - break; - } - } else { - deque.clear(); - break; - } + completableFuture.completeExceptionally(e.getCause()); + // release before handle request queue, + //in order to client reconnect infinite loop + tcLoadSemaphore.release(); + long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; + while (true) { + // prevent thread in a busy loop. + if (System.currentTimeMillis() < endTime) { + CompletableFuture future = deque.poll(); + if (future != null) { + // this means that this tc client connection connect fail + future.completeExceptionally(e); + } else { + break; } - LOG.error("Add transaction metadata store with id {} error", tcId.getId(), e); - }); - return null; - }); + } else { + deque.clear(); + break; + } + } + LOG.error("Add transaction metadata store with id {} error", tcId.getId(), e); + }); + return null; + }); } else { // only one command can open transaction metadata store, // other will be added to the deque, when the op of openTransactionMetadataStore finished @@ -448,7 +450,6 @@ private CompletableFuture endTxnInTransactionBuffer(TxnID txnID, int txnAc private static boolean isRetryableException(Throwable ex) { Throwable realCause = FutureUtil.unwrapCompletionException(ex); return (realCause instanceof TransactionMetadataStoreStateException - || realCause instanceof CoordinatorNotFoundException || realCause instanceof RequestTimeoutException || realCause instanceof ManagedLedgerException || realCause instanceof BrokerPersistenceException From 1cf93c29ef5a77243b37371ff2aaba6d5d53e73c Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 30 Jan 2023 14:38:12 +0800 Subject: [PATCH 3/9] change patch way, move recoverTracker to TransactionMetadataStoreService to decouple the mutual state dependence. --- .../TransactionMetadataStoreService.java | 114 ++++++++++-------- .../coordinator/TransactionMetadataStore.java | 6 + .../impl/MLTransactionMetadataStore.java | 7 +- 3 files changed, 73 insertions(+), 54 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index a88998b7325cd..913a65f55fca7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -40,6 +40,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import org.apache.bookkeeper.mledger.ManagedLedgerException; +import org.apache.commons.lang3.tuple.MutablePair; import org.apache.pulsar.broker.service.BrokerServiceException.ServiceUnitNotReadyException; import org.apache.pulsar.broker.transaction.exception.coordinator.TransactionCoordinatorException; import org.apache.pulsar.broker.transaction.recover.TransactionRecoverTrackerImpl; @@ -133,56 +134,61 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc return; } - openTransactionMetadataStore(tcId).thenAccept((store) -> { + openTransactionMetadataStore(tcId).thenAccept((pair) -> internalPinnedExecutor.execute(() -> { + // TransactionMetadataStore initialization need to use TransactionMetadataStore itself. + // we need to put store into stores map before handle committing and aborting transaction. + TransactionMetadataStore store = pair.getLeft(); + TransactionRecoverTracker recoverTracker = pair.getRight(); stores.put(tcId, store); LOG.info("Added new transaction meta store {}", tcId); - internalPinnedExecutor.execute(() -> { - long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; - while (true) { - // prevent thread in a busy loop. - if (System.currentTimeMillis() < endTime) { - CompletableFuture future = deque.poll(); - if (future != null) { - // complete queue request future - future.complete(null); - } else { - break; - } + recoverTracker.handleCommittingAndAbortingTransaction(); + store.setRecoverEndTime(System.currentTimeMillis()); + + long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; + while (true) { + // prevent thread in a busy loop. + if (System.currentTimeMillis() < endTime) { + CompletableFuture future = deque.poll(); + if (future != null) { + // complete queue request future + future.complete(null); } else { - deque.clear(); break; } + } else { + deque.clear(); + break; } + } - completableFuture.complete(null); - tcLoadSemaphore.release(); - }); - }).exceptionally(e -> { + completableFuture.complete(null); + tcLoadSemaphore.release(); + })).exceptionally(e -> { internalPinnedExecutor.execute(() -> { - completableFuture.completeExceptionally(e.getCause()); - // release before handle request queue, - //in order to client reconnect infinite loop - tcLoadSemaphore.release(); - long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; - while (true) { - // prevent thread in a busy loop. - if (System.currentTimeMillis() < endTime) { - CompletableFuture future = deque.poll(); - if (future != null) { - // this means that this tc client connection connect fail - future.completeExceptionally(e); - } else { - break; + completableFuture.completeExceptionally(e.getCause()); + // release before handle request queue, + //in order to client reconnect infinite loop + tcLoadSemaphore.release(); + long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; + while (true) { + // prevent thread in a busy loop. + if (System.currentTimeMillis() < endTime) { + CompletableFuture future = deque.poll(); + if (future != null) { + // this means that this tc client connection connect fail + future.completeExceptionally(e); + } else { + break; + } + } else { + deque.clear(); + break; + } } - } else { - deque.clear(); - break; - } - } - LOG.error("Add transaction metadata store with id {} error", tcId.getId(), e); - }); - return null; - }); + LOG.error("Add transaction metadata store with id {} error", tcId.getId(), e); + }); + return null; + }); } else { // only one command can open transaction metadata store, // other will be added to the deque, when the op of openTransactionMetadataStore finished @@ -202,9 +208,9 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc return completableFuture; } - public CompletableFuture openTransactionMetadataStore(TransactionCoordinatorID tcId) { - final Timer brokerClientSharedTimer = - pulsarService.getBrokerClientSharedTimer(); + public CompletableFuture> openTransactionMetadataStore(TransactionCoordinatorID tcId) { + final Timer brokerClientSharedTimer = pulsarService.getBrokerClientSharedTimer(); final ServiceConfiguration serviceConfiguration = pulsarService.getConfiguration(); final TxnLogBufferedWriterConfig txnLogBufferedWriterConfig = new TxnLogBufferedWriterConfig(); txnLogBufferedWriterConfig.setBatchEnabled(serviceConfiguration.isTransactionLogBatchedWriteEnabled()); @@ -216,15 +222,18 @@ public CompletableFuture openTransactionMetadataStore( return pulsarService.getBrokerService() .getManagedLedgerConfig(getMLTransactionLogName(tcId)).thenCompose(v -> { - TransactionTimeoutTracker timeoutTracker = timeoutTrackerFactory.newTracker(tcId); - TransactionRecoverTracker recoverTracker = - new TransactionRecoverTrackerImpl(TransactionMetadataStoreService.this, + TransactionTimeoutTracker timeoutTracker = timeoutTrackerFactory.newTracker(tcId); + TransactionRecoverTracker recoverTracker = + new TransactionRecoverTrackerImpl(TransactionMetadataStoreService.this, timeoutTracker, tcId.getId()); - return transactionMetadataStoreProvider - .openStore(tcId, pulsarService.getManagedLedgerFactory(), v, - timeoutTracker, recoverTracker, - pulsarService.getConfig().getMaxActiveTransactionsPerCoordinator(), - txnLogBufferedWriterConfig, brokerClientSharedTimer); + CompletableFuture> + completableFuture = new CompletableFuture<>(); + transactionMetadataStoreProvider.openStore(tcId, pulsarService.getManagedLedgerFactory(), v, + timeoutTracker, recoverTracker, + pulsarService.getConfig().getMaxActiveTransactionsPerCoordinator(), + txnLogBufferedWriterConfig, brokerClientSharedTimer) + .thenAccept(store -> completableFuture.complete(new MutablePair<>(store, recoverTracker))); + return completableFuture; }); } @@ -450,6 +459,7 @@ private CompletableFuture endTxnInTransactionBuffer(TxnID txnID, int txnAc private static boolean isRetryableException(Throwable ex) { Throwable realCause = FutureUtil.unwrapCompletionException(ex); return (realCause instanceof TransactionMetadataStoreStateException + || realCause instanceof CoordinatorNotFoundException || realCause instanceof RequestTimeoutException || realCause instanceof ManagedLedgerException || realCause instanceof BrokerPersistenceException diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java index 7f8280f5226b3..0d9a94df18b42 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java @@ -138,4 +138,10 @@ default long getLowWaterMark() { * @return {@link TxnMeta} the txnMetas of slow transactions */ List getSlowTransactions(long timeout); + + /** + * set recover end time. + * @param time + */ + void setRecoverEndTime(long time); } diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java index 53a515ff99164..b2e6ab82e9a2b 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java @@ -128,9 +128,7 @@ public void replayComplete() { } else { completableFuture.complete(MLTransactionMetadataStore.this); - recoverTracker.handleCommittingAndAbortingTransaction(); timeoutTracker.start(); - recoverTime.setRecoverEndTime(System.currentTimeMillis()); } } @@ -512,6 +510,11 @@ public List getSlowTransactions(long timeout) { return txnMetas; } + @Override + public void setRecoverEndTime(long time) { + recoverTime.setRecoverEndTime(time); + } + public static List txnSubscriptionToSubscription(List tnxSubscriptions) { List subscriptions = new ArrayList<>(tnxSubscriptions.size()); for (TransactionSubscription transactionSubscription : tnxSubscriptions) { From 2a8f40537705647ded4d418d7bbd97cecc5dbcd4 Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 30 Jan 2023 14:40:50 +0800 Subject: [PATCH 4/9] fix. --- .../apache/pulsar/broker/TransactionMetadataStoreService.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index 913a65f55fca7..f07593ba13ae9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -459,7 +459,6 @@ private CompletableFuture endTxnInTransactionBuffer(TxnID txnID, int txnAc private static boolean isRetryableException(Throwable ex) { Throwable realCause = FutureUtil.unwrapCompletionException(ex); return (realCause instanceof TransactionMetadataStoreStateException - || realCause instanceof CoordinatorNotFoundException || realCause instanceof RequestTimeoutException || realCause instanceof ManagedLedgerException || realCause instanceof BrokerPersistenceException From 426883fc722178e83324f2891762362c278b87bc Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 30 Jan 2023 16:28:46 +0800 Subject: [PATCH 5/9] change. --- .../TransactionMetadataStoreService.java | 134 +++++++++--------- 1 file changed, 65 insertions(+), 69 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index f07593ba13ae9..ac4691c7a0e86 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -40,7 +40,6 @@ import java.util.stream.Collectors; import java.util.stream.Stream; import org.apache.bookkeeper.mledger.ManagedLedgerException; -import org.apache.commons.lang3.tuple.MutablePair; import org.apache.pulsar.broker.service.BrokerServiceException.ServiceUnitNotReadyException; import org.apache.pulsar.broker.transaction.exception.coordinator.TransactionCoordinatorException; import org.apache.pulsar.broker.transaction.recover.TransactionRecoverTrackerImpl; @@ -134,61 +133,66 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc return; } - openTransactionMetadataStore(tcId).thenAccept((pair) -> internalPinnedExecutor.execute(() -> { - // TransactionMetadataStore initialization need to use TransactionMetadataStore itself. - // we need to put store into stores map before handle committing and aborting transaction. - TransactionMetadataStore store = pair.getLeft(); - TransactionRecoverTracker recoverTracker = pair.getRight(); - stores.put(tcId, store); - LOG.info("Added new transaction meta store {}", tcId); - recoverTracker.handleCommittingAndAbortingTransaction(); - store.setRecoverEndTime(System.currentTimeMillis()); - - long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; - while (true) { - // prevent thread in a busy loop. - if (System.currentTimeMillis() < endTime) { - CompletableFuture future = deque.poll(); - if (future != null) { - // complete queue request future - future.complete(null); - } else { - break; - } - } else { - deque.clear(); - break; - } - } - - completableFuture.complete(null); - tcLoadSemaphore.release(); - })).exceptionally(e -> { - internalPinnedExecutor.execute(() -> { - completableFuture.completeExceptionally(e.getCause()); - // release before handle request queue, - //in order to client reconnect infinite loop - tcLoadSemaphore.release(); - long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; - while (true) { - // prevent thread in a busy loop. - if (System.currentTimeMillis() < endTime) { - CompletableFuture future = deque.poll(); - if (future != null) { - // this means that this tc client connection connect fail - future.completeExceptionally(e); - } else { - break; - } + TransactionTimeoutTracker timeoutTracker = timeoutTrackerFactory.newTracker(tcId); + TransactionRecoverTracker recoverTracker = + new TransactionRecoverTrackerImpl(TransactionMetadataStoreService.this, + timeoutTracker, tcId.getId()); + openTransactionMetadataStore(tcId, timeoutTracker, recoverTracker).thenAccept( + (store) -> internalPinnedExecutor.execute(() -> { + // TransactionMetadataStore initialization + // need to use TransactionMetadataStore itself. + // we need to put store into stores map before + // handle committing and aborting transaction. + stores.put(tcId, store); + LOG.info("Added new transaction meta store {}", tcId); + recoverTracker.handleCommittingAndAbortingTransaction(); + store.setRecoverEndTime(System.currentTimeMillis()); + + long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; + while (true) { + // prevent thread in a busy loop. + if (System.currentTimeMillis() < endTime) { + CompletableFuture future = deque.poll(); + if (future != null) { + // complete queue request future + future.complete(null); } else { - deque.clear(); break; } + } else { + deque.clear(); + break; } - LOG.error("Add transaction metadata store with id {} error", tcId.getId(), e); - }); - return null; - }); + } + + completableFuture.complete(null); + tcLoadSemaphore.release(); + })).exceptionally(e -> { + internalPinnedExecutor.execute(() -> { + completableFuture.completeExceptionally(e.getCause()); + // release before handle request queue, + //in order to client reconnect infinite loop + tcLoadSemaphore.release(); + long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; + while (true) { + // prevent thread in a busy loop. + if (System.currentTimeMillis() < endTime) { + CompletableFuture future = deque.poll(); + if (future != null) { + // this means that this tc client connection connect fail + future.completeExceptionally(e); + } else { + break; + } + } else { + deque.clear(); + break; + } + } + LOG.error("Add transaction metadata store with id {} error", tcId.getId(), e); + }); + return null; + }); } else { // only one command can open transaction metadata store, // other will be added to the deque, when the op of openTransactionMetadataStore finished @@ -208,8 +212,10 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc return completableFuture; } - public CompletableFuture> openTransactionMetadataStore(TransactionCoordinatorID tcId) { + public CompletableFuture + openTransactionMetadataStore(TransactionCoordinatorID tcId, + TransactionTimeoutTracker timeoutTracker, + TransactionRecoverTracker recoverTracker) { final Timer brokerClientSharedTimer = pulsarService.getBrokerClientSharedTimer(); final ServiceConfiguration serviceConfiguration = pulsarService.getConfiguration(); final TxnLogBufferedWriterConfig txnLogBufferedWriterConfig = new TxnLogBufferedWriterConfig(); @@ -220,21 +226,11 @@ TransactionRecoverTracker>> openTransactionMetadataStore(TransactionCoordinatorI txnLogBufferedWriterConfig .setBatchedWriteMaxDelayInMillis(serviceConfiguration.getTransactionLogBatchedWriteMaxDelayInMillis()); - return pulsarService.getBrokerService() - .getManagedLedgerConfig(getMLTransactionLogName(tcId)).thenCompose(v -> { - TransactionTimeoutTracker timeoutTracker = timeoutTrackerFactory.newTracker(tcId); - TransactionRecoverTracker recoverTracker = - new TransactionRecoverTrackerImpl(TransactionMetadataStoreService.this, - timeoutTracker, tcId.getId()); - CompletableFuture> - completableFuture = new CompletableFuture<>(); - transactionMetadataStoreProvider.openStore(tcId, pulsarService.getManagedLedgerFactory(), v, - timeoutTracker, recoverTracker, - pulsarService.getConfig().getMaxActiveTransactionsPerCoordinator(), - txnLogBufferedWriterConfig, brokerClientSharedTimer) - .thenAccept(store -> completableFuture.complete(new MutablePair<>(store, recoverTracker))); - return completableFuture; - }); + return pulsarService.getBrokerService().getManagedLedgerConfig(getMLTransactionLogName(tcId)).thenCompose( + v -> transactionMetadataStoreProvider.openStore(tcId, pulsarService.getManagedLedgerFactory(), v, + timeoutTracker, recoverTracker, + pulsarService.getConfig().getMaxActiveTransactionsPerCoordinator(), txnLogBufferedWriterConfig, + brokerClientSharedTimer)); } public CompletableFuture removeTransactionMetadataStore(TransactionCoordinatorID tcId) { From e9f2f30a00c2ca83567fa5786b7d9216046cba9b Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 30 Jan 2023 17:24:33 +0800 Subject: [PATCH 6/9] fix. --- .../pulsar/broker/TransactionMetadataStoreService.java | 10 ++++++++-- .../coordinator/TransactionMetadataStore.java | 6 ------ .../coordinator/impl/MLTransactionMetadataStore.java | 5 ----- 3 files changed, 8 insertions(+), 13 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index ac4691c7a0e86..fae096426f916 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -67,6 +67,7 @@ import org.apache.pulsar.transaction.coordinator.exceptions.CoordinatorException.CoordinatorNotFoundException; import org.apache.pulsar.transaction.coordinator.exceptions.CoordinatorException.InvalidTxnStatusException; import org.apache.pulsar.transaction.coordinator.exceptions.CoordinatorException.TransactionMetadataStoreStateException; +import org.apache.pulsar.transaction.coordinator.impl.MLTransactionMetadataStore; import org.apache.pulsar.transaction.coordinator.impl.TxnLogBufferedWriterConfig; import org.apache.pulsar.transaction.coordinator.proto.TxnStatus; import org.slf4j.Logger; @@ -138,7 +139,7 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc new TransactionRecoverTrackerImpl(TransactionMetadataStoreService.this, timeoutTracker, tcId.getId()); openTransactionMetadataStore(tcId, timeoutTracker, recoverTracker).thenAccept( - (store) -> internalPinnedExecutor.execute(() -> { + store -> internalPinnedExecutor.execute(() -> { // TransactionMetadataStore initialization // need to use TransactionMetadataStore itself. // we need to put store into stores map before @@ -146,7 +147,12 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc stores.put(tcId, store); LOG.info("Added new transaction meta store {}", tcId); recoverTracker.handleCommittingAndAbortingTransaction(); - store.setRecoverEndTime(System.currentTimeMillis()); + if (store instanceof MLTransactionMetadataStore) { + MLTransactionMetadataStore mlTransactionMetadataStore = + (MLTransactionMetadataStore) store; + mlTransactionMetadataStore.recoverTime.setRecoverEndTime( + System.currentTimeMillis()); + } long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; while (true) { diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java index 0d9a94df18b42..7f8280f5226b3 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/TransactionMetadataStore.java @@ -138,10 +138,4 @@ default long getLowWaterMark() { * @return {@link TxnMeta} the txnMetas of slow transactions */ List getSlowTransactions(long timeout); - - /** - * set recover end time. - * @param time - */ - void setRecoverEndTime(long time); } diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java index b2e6ab82e9a2b..633e6b33e2137 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java @@ -510,11 +510,6 @@ public List getSlowTransactions(long timeout) { return txnMetas; } - @Override - public void setRecoverEndTime(long time) { - recoverTime.setRecoverEndTime(time); - } - public static List txnSubscriptionToSubscription(List tnxSubscriptions) { List subscriptions = new ArrayList<>(tnxSubscriptions.size()); for (TransactionSubscription transactionSubscription : tnxSubscriptions) { From 1a1e07b19848dea45776c76701f9c39f10960204 Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 30 Jan 2023 18:15:06 +0800 Subject: [PATCH 7/9] fix. --- .../pulsar/broker/TransactionMetadataStoreService.java | 6 ------ .../coordinator/impl/MLTransactionMetadataStore.java | 1 + 2 files changed, 1 insertion(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index fae096426f916..d46e663521594 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -147,12 +147,6 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc stores.put(tcId, store); LOG.info("Added new transaction meta store {}", tcId); recoverTracker.handleCommittingAndAbortingTransaction(); - if (store instanceof MLTransactionMetadataStore) { - MLTransactionMetadataStore mlTransactionMetadataStore = - (MLTransactionMetadataStore) store; - mlTransactionMetadataStore.recoverTime.setRecoverEndTime( - System.currentTimeMillis()); - } long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; while (true) { diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java index 633e6b33e2137..9a89e4d6da682 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java @@ -129,6 +129,7 @@ public void replayComplete() { } else { completableFuture.complete(MLTransactionMetadataStore.this); timeoutTracker.start(); + recoverTime.setRecoverEndTime(System.currentTimeMillis()); } } From d5199ae4df19b10b335335650dc395725fd4ba73 Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Mon, 30 Jan 2023 18:23:49 +0800 Subject: [PATCH 8/9] fix check. --- .../apache/pulsar/broker/TransactionMetadataStoreService.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index d46e663521594..e4ba73086d9dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -67,7 +67,6 @@ import org.apache.pulsar.transaction.coordinator.exceptions.CoordinatorException.CoordinatorNotFoundException; import org.apache.pulsar.transaction.coordinator.exceptions.CoordinatorException.InvalidTxnStatusException; import org.apache.pulsar.transaction.coordinator.exceptions.CoordinatorException.TransactionMetadataStoreStateException; -import org.apache.pulsar.transaction.coordinator.impl.MLTransactionMetadataStore; import org.apache.pulsar.transaction.coordinator.impl.TxnLogBufferedWriterConfig; import org.apache.pulsar.transaction.coordinator.proto.TxnStatus; import org.slf4j.Logger; From c7c2455405d7bfa6080bb371fadb7da7e12ecc42 Mon Sep 17 00:00:00 2001 From: thetumbled <843221020@qq.com> Date: Tue, 31 Jan 2023 10:46:48 +0800 Subject: [PATCH 9/9] fix. --- .../apache/pulsar/broker/TransactionMetadataStoreService.java | 1 + .../transaction/coordinator/impl/MLTransactionMetadataStore.java | 1 - 2 files changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java index e4ba73086d9dc..d0cf22a86533a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/TransactionMetadataStoreService.java @@ -146,6 +146,7 @@ public CompletableFuture handleTcClientConnect(TransactionCoordinatorID tc stores.put(tcId, store); LOG.info("Added new transaction meta store {}", tcId); recoverTracker.handleCommittingAndAbortingTransaction(); + timeoutTracker.start(); long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT; while (true) { diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java index 9a89e4d6da682..2db278f1341be 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionMetadataStore.java @@ -128,7 +128,6 @@ public void replayComplete() { } else { completableFuture.complete(MLTransactionMetadataStore.this); - timeoutTracker.start(); recoverTime.setRecoverEndTime(System.currentTimeMillis()); } }