Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -186,54 +186,68 @@ public CompletableFuture<Void> handleTcClientConnect(TransactionCoordinatorID tc
return;
}

openTransactionMetadataStore(tcId).thenAccept((store) -> internalPinnedExecutor.execute(() -> {
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<Void> future = deque.poll();
if (future != null) {
// complete queue request future
future.complete(null);
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();
timeoutTracker.start();

long endTime = System.currentTimeMillis() + HANDLE_PENDING_CONNECT_TIME_OUT;
while (true) {

// prevent thread in a busy loop.
if (System.currentTimeMillis() < endTime) {
CompletableFuture<Void> 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<Void> 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);
});

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<Void> 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;
});
return null;
});
} else {
// only one command can open transaction metadata store,
// other will be added to the deque, when the op of openTransactionMetadataStore finished
Expand All @@ -253,16 +267,14 @@ public CompletableFuture<Void> handleTcClientConnect(TransactionCoordinatorID tc
return completableFuture;
}

public CompletableFuture<TransactionMetadataStore> openTransactionMetadataStore(TransactionCoordinatorID tcId) {
public CompletableFuture<TransactionMetadataStore>
openTransactionMetadataStore(TransactionCoordinatorID tcId,
TransactionTimeoutTracker timeoutTracker,
TransactionRecoverTracker recoverTracker) {
return pulsarService.getBrokerService()
.getManagedLedgerConfig(getMLTransactionLogName(tcId)).thenCompose(v -> {
TransactionTimeoutTracker timeoutTracker = timeoutTrackerFactory.newTracker(tcId);
TransactionRecoverTracker recoverTracker =
new TransactionRecoverTrackerImpl(TransactionMetadataStoreService.this,
timeoutTracker, tcId.getId());
return transactionMetadataStoreProvider
.openStore(tcId, pulsarService.getManagedLedgerFactory(), v,
timeoutTracker, recoverTracker);
return transactionMetadataStoreProvider.openStore(tcId,
pulsarService.getManagedLedgerFactory(), v, timeoutTracker, recoverTracker);
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -120,8 +120,6 @@ public void replayComplete() {
+ tcID.toString() + " change state to Ready error when init it"));

} else {
recoverTracker.handleCommittingAndAbortingTransaction();
timeoutTracker.start();
completableFuture.complete(MLTransactionMetadataStore.this);
}
}
Expand Down