From 542cc5a9b34a8f99433af5a8465915fb2c4e248b Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 8 Mar 2023 16:47:37 +0800 Subject: [PATCH 01/11] [fix][metadata] Fix notification thread block causes double owner. --- .../coordination/impl/LockManagerImpl.java | 40 +++++++--- .../coordination/impl/ResourceLockImpl.java | 73 ++++++++++++------- 2 files changed, 76 insertions(+), 37 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java index 097f15af27677..58be07f619835 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java @@ -19,6 +19,7 @@ package org.apache.pulsar.metadata.coordination.impl; import com.fasterxml.jackson.databind.type.TypeFactory; +import java.time.Duration; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -29,6 +30,9 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.util.SafeRunnable; @@ -48,7 +52,7 @@ @Slf4j class LockManagerImpl implements LockManager { - + private static final Duration REVALIDATE_TIMEOUT = Duration.ofSeconds(30); private final Map> locks = new ConcurrentHashMap<>(); private final MetadataStoreExtended store; private final MetadataCache cache; @@ -118,28 +122,42 @@ public CompletableFuture> acquireLock(String path, T value) { private void handleSessionEvent(SessionEvent se) { // We want to make sure we're processing one event at a time and that we're done with one event before going // for the next one. - executor.execute(SafeRunnable.safeRun(() -> { - List> futures = new ArrayList<>(); - + final SafeRunnable task = SafeRunnable.safeRun(() -> { + final List> futures = new ArrayList<>(); if (se == SessionEvent.SessionReestablished) { log.info("Metadata store session has been re-established. Revalidating all the existing locks."); for (ResourceLockImpl lock : locks.values()) { - futures.add(lock.revalidate(lock.getValue(), true)); + futures.add(lock.revalidateOnce(lock.getValue())); } - } else if (se == SessionEvent.Reconnected) { log.info("Metadata store connection has been re-established. Revalidating locks that were pending."); for (ResourceLockImpl lock : locks.values()) { futures.add(lock.revalidateIfNeededAfterReconnection()); } } - try { - FutureUtil.waitForAll(futures).get(); - } catch (ExecutionException | InterruptedException e) { - log.warn("Failure when processing session event", e); + FutureUtil.waitForAll(futures).get(REVALIDATE_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); + } catch (ExecutionException ex) { + log.warn("Got exception when execute revalidate ", ex); + } catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + log.warn("Got thread interrupted exception when execute revalidate."); + } catch (TimeoutException ex) { + log.warn("Got timeout exception when execute revalidate"); + for (final CompletableFuture future : futures) { + if (!future.isDone()) { + if(!future.cancel(true)) { + log.warn("Failed to cancel the revalidation future {}", future); + } + } + } } - })); + }); + try { + executor.execute(task); + } catch (RejectedExecutionException ex) { + log.warn("Session events cannot be executed because the executor has been closed."); + } } private void handleDataNotification(Notification n) { diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java index dea9aa1acb90f..090becd2bb26b 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java @@ -20,6 +20,7 @@ import java.util.EnumSet; import java.util.Optional; +import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.concurrent.FutureUtils; @@ -32,6 +33,7 @@ import org.apache.pulsar.metadata.api.coordination.ResourceLock; import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import javax.annotation.Nonnull; @Slf4j public class ResourceLockImpl implements ResourceLock { @@ -128,7 +130,7 @@ synchronized CompletableFuture acquire(T newValue) { .thenRun(() -> result.complete(null)) .exceptionally(ex -> { if (ex.getCause() instanceof LockBusyException) { - revalidate(newValue, false) + revalidate(newValue) .thenAccept(__ -> result.complete(null)) .exceptionally(ex1 -> { result.completeExceptionally(ex1); @@ -185,21 +187,29 @@ synchronized void lockWasInvalidated() { } log.info("Lock on resource {} was invalidated", path); - revalidate(value, true) - .thenRun(() -> log.info("Successfully revalidated the lock on {}", path)); + revalidateOnce(value); } synchronized CompletableFuture revalidateIfNeededAfterReconnection() { if (revalidateAfterReconnection) { revalidateAfterReconnection = false; log.warn("Revalidate lock at {} after reconnection", path); - return revalidate(value, true); + return revalidateOnce(value); } else { return CompletableFuture.completedFuture(null); } } - synchronized CompletableFuture revalidate(T newValue, boolean revalidateAfterReconnection) { + /** + * Revalidate the distributed lock if it is not released. + * This method is thread-safe and it will perform multiple re-validation operations in turn. + * @param newValue the lock value + */ + synchronized @Nonnull CompletableFuture revalidate(@Nonnull T newValue) { + if (state == State.Released) { + // We don't need to revalidate the released lock since the expired future has been executed. + return CompletableFuture.completedFuture(null); + } if (revalidateFuture == null || revalidateFuture.isDone()) { revalidateFuture = doRevalidate(newValue); } else { @@ -217,30 +227,41 @@ synchronized CompletableFuture revalidate(T newValue, boolean revalidateAf }); revalidateFuture = newFuture; } - revalidateFuture.exceptionally(ex -> { - synchronized (ResourceLockImpl.this) { - Throwable realCause = FutureUtil.unwrapCompletionException(ex); - if (!revalidateAfterReconnection || realCause instanceof BadVersionException - || realCause instanceof LockBusyException) { - log.warn("Failed to revalidate the lock at {}. Marked as expired. {}", - path, realCause.getMessage()); - state = State.Released; - expiredFuture.complete(null); - } else { - // We failed to revalidate the lock due to connectivity issue - // Continue assuming we hold the lock, until we can revalidate it, either - // on Reconnected or SessionReestablished events. - ResourceLockImpl.this.revalidateAfterReconnection = true; - log.warn("Failed to revalidate the lock at {}. Retrying later on reconnection {}", path, - realCause.getMessage()); - } - } - return null; - }); return revalidateFuture; } - private synchronized CompletableFuture doRevalidate(T newValue) { + /** + * This method will auto mark the lock is released if revalidation operation got one of #{@code } + * @param newValue the lock value + */ + @Nonnull CompletableFuture revalidateOnce(@Nonnull T newValue) { + return revalidate(newValue) + .thenRun(() -> log.info("Successfully revalidated once the lock on {}", path)) + .exceptionally(ex -> { + synchronized (ResourceLockImpl.this) { + Throwable realCause = FutureUtil.unwrapCompletionException(ex); + if (realCause instanceof BadVersionException || realCause instanceof LockBusyException + // If the revalidation future is cancelled, + // we can assume the invoker will give up this lock in memory. + || realCause instanceof CancellationException) { + log.warn("Failed to revalidate the lock at {}. Marked as expired. {}", + path, realCause.getMessage()); + state = State.Released; + expiredFuture.complete(null); + } else { + // We failed to revalidate the lock due to connectivity issue + // Continue assuming we hold the lock, until we can revalidate it, either + // on Reconnected or SessionReestablished events. + ResourceLockImpl.this.revalidateAfterReconnection = true; + log.warn("Failed to revalidate the lock at {}. Retrying later on reconnection {}", path, + realCause.getMessage()); + } + } + return null; + }); + } + + private synchronized @Nonnull CompletableFuture doRevalidate(@Nonnull T newValue) { if (log.isDebugEnabled()) { log.debug("doRevalidate with newValue={}, version={}", newValue, version); } From 7adf86f2cb644b34012549700ccf4522747dbe22 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 8 Mar 2023 17:42:45 +0800 Subject: [PATCH 02/11] Fix the doc --- .../coordination/impl/LockManagerImpl.java | 2 +- .../coordination/impl/ResourceLockImpl.java | 15 +++++++++++---- 2 files changed, 12 insertions(+), 5 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java index 58be07f619835..55eb3d0bb06f7 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java @@ -127,7 +127,7 @@ private void handleSessionEvent(SessionEvent se) { if (se == SessionEvent.SessionReestablished) { log.info("Metadata store session has been re-established. Revalidating all the existing locks."); for (ResourceLockImpl lock : locks.values()) { - futures.add(lock.revalidateOnce(lock.getValue())); + futures.add(lock.silentRevalidateOnce(lock.getValue())); } } else if (se == SessionEvent.Reconnected) { log.info("Metadata store connection has been re-established. Revalidating locks that were pending."); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java index 090becd2bb26b..15a2054bd6a12 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java @@ -187,14 +187,14 @@ synchronized void lockWasInvalidated() { } log.info("Lock on resource {} was invalidated", path); - revalidateOnce(value); + silentRevalidateOnce(value); } synchronized CompletableFuture revalidateIfNeededAfterReconnection() { if (revalidateAfterReconnection) { revalidateAfterReconnection = false; log.warn("Revalidate lock at {} after reconnection", path); - return revalidateOnce(value); + return silentRevalidateOnce(value); } else { return CompletableFuture.completedFuture(null); } @@ -231,10 +231,17 @@ synchronized CompletableFuture revalidateIfNeededAfterReconnection() { } /** - * This method will auto mark the lock is released if revalidation operation got one of #{@code } + * This method designed for background notification usage,it will auto mark the lock is released if revalidation + * operation got one of exceptions as follows: + * - LockBusyException + * - BadVersionException + * - CancellationException + * * @param newValue the lock value + * @return The revalidation future #Notice: It will not return any useful result, + * the caller needs to re-check the lock state after silent revalidation once. */ - @Nonnull CompletableFuture revalidateOnce(@Nonnull T newValue) { + @Nonnull CompletableFuture silentRevalidateOnce(@Nonnull T newValue) { return revalidate(newValue) .thenRun(() -> log.info("Successfully revalidated once the lock on {}", path)) .exceptionally(ex -> { From 1e4d63883dd6ef5078f87bda12a4eecfef602c67 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 8 Mar 2023 17:59:20 +0800 Subject: [PATCH 03/11] Fix checkstyle --- .../pulsar/metadata/coordination/impl/LockManagerImpl.java | 2 +- .../pulsar/metadata/coordination/impl/ResourceLockImpl.java | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java index 55eb3d0bb06f7..a8e5af4f3a13d 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java @@ -146,7 +146,7 @@ private void handleSessionEvent(SessionEvent se) { log.warn("Got timeout exception when execute revalidate"); for (final CompletableFuture future : futures) { if (!future.isDone()) { - if(!future.cancel(true)) { + if (!future.cancel(true)) { log.warn("Failed to cancel the revalidation future {}", future); } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java index 15a2054bd6a12..ecaa91e9f6078 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java @@ -22,6 +22,7 @@ import java.util.Optional; import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; +import javax.annotation.Nonnull; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.concurrent.FutureUtils; import org.apache.pulsar.common.util.FutureUtil; @@ -33,7 +34,6 @@ import org.apache.pulsar.metadata.api.coordination.ResourceLock; import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; -import javax.annotation.Nonnull; @Slf4j public class ResourceLockImpl implements ResourceLock { @@ -231,8 +231,8 @@ synchronized CompletableFuture revalidateIfNeededAfterReconnection() { } /** - * This method designed for background notification usage,it will auto mark the lock is released if revalidation - * operation got one of exceptions as follows: + * This method designed for background notification usage. + * It will auto mark the lock is released if revalidation operation got one of exceptions as follows: * - LockBusyException * - BadVersionException * - CancellationException From 1a6082f6b54d3899213fe8cf9ee54baf4794d3f5 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 9 Mar 2023 23:31:00 +0800 Subject: [PATCH 04/11] [fix][meta] Fix busy-loop causes watch can't acquire lock. --- .../apache/pulsar/metadata/impl/ZKSessionWatcher.java | 2 +- .../java/org/apache/pulsar/metadata/ZKSessionTest.java | 10 +++++----- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java index a0247e2231949..f2f1cfb484353 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java @@ -66,7 +66,7 @@ public ZKSessionWatcher(ZooKeeper zk, Consumer sessionListener) { this.scheduler = Executors .newSingleThreadScheduledExecutor(new DefaultThreadFactory("metadata-store-zk-session-watcher")); this.task = - scheduler.scheduleAtFixedRate(catchingAndLoggingThrowables(this::checkConnectionStatus), tickTimeMillis, + scheduler.scheduleWithFixedDelay(catchingAndLoggingThrowables(this::checkConnectionStatus), tickTimeMillis, tickTimeMillis, TimeUnit.MILLISECONDS); this.currentStatus = SessionEvent.SessionReestablished; diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java index 1757a6b991c0f..4de9904073a02 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java @@ -40,7 +40,7 @@ import org.awaitility.Awaitility; import org.testng.annotations.Test; -@Test(groups = "quarantine") +@Test public class ZKSessionTest extends BaseMetadataStoreTest { @Test @@ -139,20 +139,20 @@ public void testReacquireLocksAfterSessionLost() throws Exception { @Test public void testReacquireLeadershipAfterSessionLost() throws Exception { // --- init - @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(zks.getConnectionString(), MetadataStoreConfig.builder() .sessionTimeoutMillis(2_000) .build()); BlockingQueue sessionEvents = new LinkedBlockingQueue<>(); - store.registerSessionListener(sessionEvents::add); + store.registerSessionListener(event -> { + System.out.println("event:" + event); + sessionEvents.add(event); + }); BlockingQueue leaderElectionEvents = new LinkedBlockingQueue<>(); String path = newKey(); - @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); - @Cleanup LeaderElection le1 = coordinationService.getLeaderElection(String.class, path, leaderElectionEvents::add); // --- test manual elect From 3ff104122691247ed242b975170676fb75885ff2 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 9 Mar 2023 23:33:14 +0800 Subject: [PATCH 05/11] Revert "Fix checkstyle" This reverts commit 1e4d63883dd6ef5078f87bda12a4eecfef602c67. --- .../pulsar/metadata/coordination/impl/LockManagerImpl.java | 2 +- .../pulsar/metadata/coordination/impl/ResourceLockImpl.java | 6 +++--- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java index a8e5af4f3a13d..55eb3d0bb06f7 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java @@ -146,7 +146,7 @@ private void handleSessionEvent(SessionEvent se) { log.warn("Got timeout exception when execute revalidate"); for (final CompletableFuture future : futures) { if (!future.isDone()) { - if (!future.cancel(true)) { + if(!future.cancel(true)) { log.warn("Failed to cancel the revalidation future {}", future); } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java index ecaa91e9f6078..15a2054bd6a12 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java @@ -22,7 +22,6 @@ import java.util.Optional; import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; -import javax.annotation.Nonnull; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.concurrent.FutureUtils; import org.apache.pulsar.common.util.FutureUtil; @@ -34,6 +33,7 @@ import org.apache.pulsar.metadata.api.coordination.ResourceLock; import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import javax.annotation.Nonnull; @Slf4j public class ResourceLockImpl implements ResourceLock { @@ -231,8 +231,8 @@ synchronized CompletableFuture revalidateIfNeededAfterReconnection() { } /** - * This method designed for background notification usage. - * It will auto mark the lock is released if revalidation operation got one of exceptions as follows: + * This method designed for background notification usage,it will auto mark the lock is released if revalidation + * operation got one of exceptions as follows: * - LockBusyException * - BadVersionException * - CancellationException From ccd8eff42c34de3eafd7db21c3cbf9a8a8cf1b45 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 9 Mar 2023 23:33:14 +0800 Subject: [PATCH 06/11] Revert "Fix the doc" This reverts commit 7adf86f2cb644b34012549700ccf4522747dbe22. --- .../coordination/impl/LockManagerImpl.java | 2 +- .../coordination/impl/ResourceLockImpl.java | 15 ++++----------- 2 files changed, 5 insertions(+), 12 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java index 55eb3d0bb06f7..58be07f619835 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java @@ -127,7 +127,7 @@ private void handleSessionEvent(SessionEvent se) { if (se == SessionEvent.SessionReestablished) { log.info("Metadata store session has been re-established. Revalidating all the existing locks."); for (ResourceLockImpl lock : locks.values()) { - futures.add(lock.silentRevalidateOnce(lock.getValue())); + futures.add(lock.revalidateOnce(lock.getValue())); } } else if (se == SessionEvent.Reconnected) { log.info("Metadata store connection has been re-established. Revalidating locks that were pending."); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java index 15a2054bd6a12..090becd2bb26b 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java @@ -187,14 +187,14 @@ synchronized void lockWasInvalidated() { } log.info("Lock on resource {} was invalidated", path); - silentRevalidateOnce(value); + revalidateOnce(value); } synchronized CompletableFuture revalidateIfNeededAfterReconnection() { if (revalidateAfterReconnection) { revalidateAfterReconnection = false; log.warn("Revalidate lock at {} after reconnection", path); - return silentRevalidateOnce(value); + return revalidateOnce(value); } else { return CompletableFuture.completedFuture(null); } @@ -231,17 +231,10 @@ synchronized CompletableFuture revalidateIfNeededAfterReconnection() { } /** - * This method designed for background notification usage,it will auto mark the lock is released if revalidation - * operation got one of exceptions as follows: - * - LockBusyException - * - BadVersionException - * - CancellationException - * + * This method will auto mark the lock is released if revalidation operation got one of #{@code } * @param newValue the lock value - * @return The revalidation future #Notice: It will not return any useful result, - * the caller needs to re-check the lock state after silent revalidation once. */ - @Nonnull CompletableFuture silentRevalidateOnce(@Nonnull T newValue) { + @Nonnull CompletableFuture revalidateOnce(@Nonnull T newValue) { return revalidate(newValue) .thenRun(() -> log.info("Successfully revalidated once the lock on {}", path)) .exceptionally(ex -> { From 7127bf3a9ca1c9fcbccf08d6ee142e012e701957 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 9 Mar 2023 23:33:15 +0800 Subject: [PATCH 07/11] Revert "[fix][metadata] Fix notification thread block causes double owner." This reverts commit 542cc5a9b34a8f99433af5a8465915fb2c4e248b. --- .../coordination/impl/LockManagerImpl.java | 40 +++------- .../coordination/impl/ResourceLockImpl.java | 73 +++++++------------ 2 files changed, 37 insertions(+), 76 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java index 58be07f619835..097f15af27677 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/LockManagerImpl.java @@ -19,7 +19,6 @@ package org.apache.pulsar.metadata.coordination.impl; import com.fasterxml.jackson.databind.type.TypeFactory; -import java.time.Duration; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -30,9 +29,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; -import java.util.concurrent.RejectedExecutionException; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.util.SafeRunnable; @@ -52,7 +48,7 @@ @Slf4j class LockManagerImpl implements LockManager { - private static final Duration REVALIDATE_TIMEOUT = Duration.ofSeconds(30); + private final Map> locks = new ConcurrentHashMap<>(); private final MetadataStoreExtended store; private final MetadataCache cache; @@ -122,42 +118,28 @@ public CompletableFuture> acquireLock(String path, T value) { private void handleSessionEvent(SessionEvent se) { // We want to make sure we're processing one event at a time and that we're done with one event before going // for the next one. - final SafeRunnable task = SafeRunnable.safeRun(() -> { - final List> futures = new ArrayList<>(); + executor.execute(SafeRunnable.safeRun(() -> { + List> futures = new ArrayList<>(); + if (se == SessionEvent.SessionReestablished) { log.info("Metadata store session has been re-established. Revalidating all the existing locks."); for (ResourceLockImpl lock : locks.values()) { - futures.add(lock.revalidateOnce(lock.getValue())); + futures.add(lock.revalidate(lock.getValue(), true)); } + } else if (se == SessionEvent.Reconnected) { log.info("Metadata store connection has been re-established. Revalidating locks that were pending."); for (ResourceLockImpl lock : locks.values()) { futures.add(lock.revalidateIfNeededAfterReconnection()); } } + try { - FutureUtil.waitForAll(futures).get(REVALIDATE_TIMEOUT.toMillis(), TimeUnit.MILLISECONDS); - } catch (ExecutionException ex) { - log.warn("Got exception when execute revalidate ", ex); - } catch (InterruptedException ex) { - Thread.currentThread().interrupt(); - log.warn("Got thread interrupted exception when execute revalidate."); - } catch (TimeoutException ex) { - log.warn("Got timeout exception when execute revalidate"); - for (final CompletableFuture future : futures) { - if (!future.isDone()) { - if(!future.cancel(true)) { - log.warn("Failed to cancel the revalidation future {}", future); - } - } - } + FutureUtil.waitForAll(futures).get(); + } catch (ExecutionException | InterruptedException e) { + log.warn("Failure when processing session event", e); } - }); - try { - executor.execute(task); - } catch (RejectedExecutionException ex) { - log.warn("Session events cannot be executed because the executor has been closed."); - } + })); } private void handleDataNotification(Notification n) { diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java index 090becd2bb26b..dea9aa1acb90f 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/coordination/impl/ResourceLockImpl.java @@ -20,7 +20,6 @@ import java.util.EnumSet; import java.util.Optional; -import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.concurrent.FutureUtils; @@ -33,7 +32,6 @@ import org.apache.pulsar.metadata.api.coordination.ResourceLock; import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; -import javax.annotation.Nonnull; @Slf4j public class ResourceLockImpl implements ResourceLock { @@ -130,7 +128,7 @@ synchronized CompletableFuture acquire(T newValue) { .thenRun(() -> result.complete(null)) .exceptionally(ex -> { if (ex.getCause() instanceof LockBusyException) { - revalidate(newValue) + revalidate(newValue, false) .thenAccept(__ -> result.complete(null)) .exceptionally(ex1 -> { result.completeExceptionally(ex1); @@ -187,29 +185,21 @@ synchronized void lockWasInvalidated() { } log.info("Lock on resource {} was invalidated", path); - revalidateOnce(value); + revalidate(value, true) + .thenRun(() -> log.info("Successfully revalidated the lock on {}", path)); } synchronized CompletableFuture revalidateIfNeededAfterReconnection() { if (revalidateAfterReconnection) { revalidateAfterReconnection = false; log.warn("Revalidate lock at {} after reconnection", path); - return revalidateOnce(value); + return revalidate(value, true); } else { return CompletableFuture.completedFuture(null); } } - /** - * Revalidate the distributed lock if it is not released. - * This method is thread-safe and it will perform multiple re-validation operations in turn. - * @param newValue the lock value - */ - synchronized @Nonnull CompletableFuture revalidate(@Nonnull T newValue) { - if (state == State.Released) { - // We don't need to revalidate the released lock since the expired future has been executed. - return CompletableFuture.completedFuture(null); - } + synchronized CompletableFuture revalidate(T newValue, boolean revalidateAfterReconnection) { if (revalidateFuture == null || revalidateFuture.isDone()) { revalidateFuture = doRevalidate(newValue); } else { @@ -227,41 +217,30 @@ synchronized CompletableFuture revalidateIfNeededAfterReconnection() { }); revalidateFuture = newFuture; } + revalidateFuture.exceptionally(ex -> { + synchronized (ResourceLockImpl.this) { + Throwable realCause = FutureUtil.unwrapCompletionException(ex); + if (!revalidateAfterReconnection || realCause instanceof BadVersionException + || realCause instanceof LockBusyException) { + log.warn("Failed to revalidate the lock at {}. Marked as expired. {}", + path, realCause.getMessage()); + state = State.Released; + expiredFuture.complete(null); + } else { + // We failed to revalidate the lock due to connectivity issue + // Continue assuming we hold the lock, until we can revalidate it, either + // on Reconnected or SessionReestablished events. + ResourceLockImpl.this.revalidateAfterReconnection = true; + log.warn("Failed to revalidate the lock at {}. Retrying later on reconnection {}", path, + realCause.getMessage()); + } + } + return null; + }); return revalidateFuture; } - /** - * This method will auto mark the lock is released if revalidation operation got one of #{@code } - * @param newValue the lock value - */ - @Nonnull CompletableFuture revalidateOnce(@Nonnull T newValue) { - return revalidate(newValue) - .thenRun(() -> log.info("Successfully revalidated once the lock on {}", path)) - .exceptionally(ex -> { - synchronized (ResourceLockImpl.this) { - Throwable realCause = FutureUtil.unwrapCompletionException(ex); - if (realCause instanceof BadVersionException || realCause instanceof LockBusyException - // If the revalidation future is cancelled, - // we can assume the invoker will give up this lock in memory. - || realCause instanceof CancellationException) { - log.warn("Failed to revalidate the lock at {}. Marked as expired. {}", - path, realCause.getMessage()); - state = State.Released; - expiredFuture.complete(null); - } else { - // We failed to revalidate the lock due to connectivity issue - // Continue assuming we hold the lock, until we can revalidate it, either - // on Reconnected or SessionReestablished events. - ResourceLockImpl.this.revalidateAfterReconnection = true; - log.warn("Failed to revalidate the lock at {}. Retrying later on reconnection {}", path, - realCause.getMessage()); - } - } - return null; - }); - } - - private synchronized @Nonnull CompletableFuture doRevalidate(@Nonnull T newValue) { + private synchronized CompletableFuture doRevalidate(T newValue) { if (log.isDebugEnabled()) { log.debug("doRevalidate with newValue={}, version={}", newValue, version); } From 0de94b7ec545fc0a59706b04e1ba867d25709900 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 9 Mar 2023 23:33:51 +0800 Subject: [PATCH 08/11] Revert some useless code --- .../java/org/apache/pulsar/metadata/ZKSessionTest.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java index 4de9904073a02..36cb0f132ba58 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java @@ -139,20 +139,20 @@ public void testReacquireLocksAfterSessionLost() throws Exception { @Test public void testReacquireLeadershipAfterSessionLost() throws Exception { // --- init + @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(zks.getConnectionString(), MetadataStoreConfig.builder() .sessionTimeoutMillis(2_000) .build()); BlockingQueue sessionEvents = new LinkedBlockingQueue<>(); - store.registerSessionListener(event -> { - System.out.println("event:" + event); - sessionEvents.add(event); - }); + store.registerSessionListener(sessionEvents::add); BlockingQueue leaderElectionEvents = new LinkedBlockingQueue<>(); String path = newKey(); + @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); + @Cleanup LeaderElection le1 = coordinationService.getLeaderElection(String.class, path, leaderElectionEvents::add); // --- test manual elect From 8f95e6686bdb4420f1c12649110d4f4b64095ab5 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 10 Mar 2023 00:06:04 +0800 Subject: [PATCH 09/11] Fix checlstyle --- .../org/apache/pulsar/metadata/impl/ZKSessionWatcher.java | 3 ++- .../java/org/apache/pulsar/metadata/ZKSessionTest.java | 7 +++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java index f2f1cfb484353..1ce01f57d4fbe 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKSessionWatcher.java @@ -66,7 +66,8 @@ public ZKSessionWatcher(ZooKeeper zk, Consumer sessionListener) { this.scheduler = Executors .newSingleThreadScheduledExecutor(new DefaultThreadFactory("metadata-store-zk-session-watcher")); this.task = - scheduler.scheduleWithFixedDelay(catchingAndLoggingThrowables(this::checkConnectionStatus), tickTimeMillis, + scheduler.scheduleWithFixedDelay( + catchingAndLoggingThrowables(this::checkConnectionStatus), tickTimeMillis, tickTimeMillis, TimeUnit.MILLISECONDS); this.currentStatus = SessionEvent.SessionReestablished; diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java index 36cb0f132ba58..7e923a5f90626 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java @@ -139,10 +139,9 @@ public void testReacquireLocksAfterSessionLost() throws Exception { @Test public void testReacquireLeadershipAfterSessionLost() throws Exception { // --- init - @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(zks.getConnectionString(), MetadataStoreConfig.builder() - .sessionTimeoutMillis(2_000) + .sessionTimeoutMillis(30000) .build()); BlockingQueue sessionEvents = new LinkedBlockingQueue<>(); @@ -150,9 +149,7 @@ public void testReacquireLeadershipAfterSessionLost() throws Exception { BlockingQueue leaderElectionEvents = new LinkedBlockingQueue<>(); String path = newKey(); - @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); - @Cleanup LeaderElection le1 = coordinationService.getLeaderElection(String.class, path, leaderElectionEvents::add); // --- test manual elect @@ -179,5 +176,7 @@ public void testReacquireLeadershipAfterSessionLost() throws Exception { Awaitility.await().atMost(Duration.ofSeconds(15)) .untilAsserted(()-> assertEquals(le1.getState(),LeaderElectionState.Leading)); assertTrue(store.get(path).join().isPresent()); + + Thread.sleep(1000000); } } From 827d62565523f3ce5f6508b72ed49ad12b522892 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 10 Mar 2023 00:08:59 +0800 Subject: [PATCH 10/11] Revert useless changes --- .../test/java/org/apache/pulsar/metadata/ZKSessionTest.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java index 7e923a5f90626..aa9ee99eb8ceb 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java @@ -141,7 +141,7 @@ public void testReacquireLeadershipAfterSessionLost() throws Exception { // --- init MetadataStoreExtended store = MetadataStoreExtended.create(zks.getConnectionString(), MetadataStoreConfig.builder() - .sessionTimeoutMillis(30000) + .sessionTimeoutMillis(2_000) .build()); BlockingQueue sessionEvents = new LinkedBlockingQueue<>(); @@ -176,7 +176,5 @@ public void testReacquireLeadershipAfterSessionLost() throws Exception { Awaitility.await().atMost(Duration.ofSeconds(15)) .untilAsserted(()-> assertEquals(le1.getState(),LeaderElectionState.Leading)); assertTrue(store.get(path).join().isPresent()); - - Thread.sleep(1000000); } } From f250d100c1957b5e452faeb9e3904a7867e230af Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 10 Mar 2023 00:09:44 +0800 Subject: [PATCH 11/11] Revert useless code --- .../test/java/org/apache/pulsar/metadata/ZKSessionTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java index aa9ee99eb8ceb..36cb0f132ba58 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/ZKSessionTest.java @@ -139,6 +139,7 @@ public void testReacquireLocksAfterSessionLost() throws Exception { @Test public void testReacquireLeadershipAfterSessionLost() throws Exception { // --- init + @Cleanup MetadataStoreExtended store = MetadataStoreExtended.create(zks.getConnectionString(), MetadataStoreConfig.builder() .sessionTimeoutMillis(2_000) @@ -149,7 +150,9 @@ public void testReacquireLeadershipAfterSessionLost() throws Exception { BlockingQueue leaderElectionEvents = new LinkedBlockingQueue<>(); String path = newKey(); + @Cleanup CoordinationService coordinationService = new CoordinationServiceImpl(store); + @Cleanup LeaderElection le1 = coordinationService.getLeaderElection(String.class, path, leaderElectionEvents::add); // --- test manual elect