From 27ddd227704a74d4e2f61d72bfbc2e688898e18a Mon Sep 17 00:00:00 2001 From: lordcheng10 Date: Thu, 5 Jan 2023 09:51:42 +0800 Subject: [PATCH 1/4] chery pick 17700 to 2.8 branch --- .../coordination/impl/LockManagerImpl.java | 4 +- .../coordination/impl/ResourceLockImpl.java | 53 +++++++++++-------- 2 files changed, 35 insertions(+), 22 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 0e3754ac29392..f7427d6eaae48 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 @@ -106,7 +106,9 @@ public CompletableFuture> acquireLock(String path, T value) { private void handleSessionEvent(SessionEvent se) { if (se == SessionEvent.SessionReestablished) { log.info("Metadata store session has been re-established. Revalidating all the existing locks."); - locks.values().forEach(ResourceLockImpl::revalidate); + for (ResourceLockImpl lock : locks.values()) { + lock.revalidate(true); + } } else if (se == SessionEvent.Reconnected) { log.info("Metadata store connection has been re-established. Revalidating locks that were pending."); locks.values().forEach(ResourceLockImpl::revalidateIfNeededAfterReconnection); 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 55b420edff38e..527ecf3db8239 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 @@ -24,6 +24,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.common.concurrent.FutureUtils; +import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.metadata.api.GetResult; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.MetadataStoreException.BadVersionException; @@ -141,7 +142,7 @@ synchronized CompletableFuture acquire() { .thenRun(() -> result.complete(null)) .exceptionally(ex -> { if (ex.getCause() instanceof LockBusyException) { - revalidate() + revalidate(false) .thenAccept(__ -> result.complete(null)) .exceptionally(ex1 -> { result.completeExceptionally(ex1); @@ -194,35 +195,19 @@ synchronized void lockWasInvalidated() { } log.info("Lock on resource {} was invalidated", path); - revalidate() - .thenRun(() -> log.info("Successfully revalidated the lock on {}", path)) - .exceptionally(ex -> { - synchronized (ResourceLockImpl.this) { - if (ex.getCause() instanceof BadVersionException) { - log.warn("Failed to revalidate the lock at {}. Marked as expired", path); - 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. - log.warn("Failed to revalidate the lock at {}. Retrying later on reconnection", path, - ex.getCause().getMessage()); - } - } - return null; - }); + revalidate(true) + .thenRun(() -> log.info("Successfully revalidated the lock on {}", path)); } synchronized void revalidateIfNeededAfterReconnection() { if (revalidateAfterReconnection) { revalidateAfterReconnection = false; log.warn("Revalidate lock at {} after reconnection", path); - revalidate(); + revalidate(true); } } - synchronized CompletableFuture revalidate() { + private synchronized CompletableFuture doRevalidate() { return store.get(path) .thenCompose(optGetResult -> { if (!optGetResult.isPresent()) { @@ -279,4 +264,30 @@ synchronized CompletableFuture revalidate() { } }); } + + synchronized CompletableFuture revalidate(boolean revalidateAfterReconnection) { + CompletableFuture revalidateFuture = doRevalidate(); + 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; + } } From 519857a04f560d0dec047267723c7db16e156027 Mon Sep 17 00:00:00 2001 From: lordcheng10 Date: Thu, 5 Jan 2023 10:16:05 +0800 Subject: [PATCH 2/4] add test testCleanUpStateWhenRevalidationGotLockBusy --- .../pulsar/metadata/LockManagerTest.java | 66 +++++++++++++++++-- 1 file changed, 62 insertions(+), 4 deletions(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java index e256950d6eaa3..0df5cadc679d6 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java @@ -18,10 +18,6 @@ */ package org.apache.pulsar.metadata; -import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertFalse; -import static org.testng.Assert.fail; - import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Arrays; @@ -31,6 +27,8 @@ import java.util.concurrent.CompletionException; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Supplier; import lombok.Cleanup; @@ -45,8 +43,12 @@ import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; import org.apache.pulsar.metadata.coordination.impl.CoordinationServiceImpl; +import org.awaitility.Awaitility; import org.testng.annotations.Test; +import static org.testng.Assert.*; +import static org.testng.Assert.assertTrue; + public class LockManagerTest extends BaseMetadataStoreTest { @Test(dataProvider = "impl") @@ -228,4 +230,60 @@ public void revalidateLockOnDifferentSession(String provider, String url) throws assertEquals(new String(store1.get(path2).join().get().getValue()), "\"value-1\""); } + + @Test(dataProvider = "impl") + public void testCleanUpStateWhenRevalidationGotLockBusy(String provider, Supplier urlSupplier) + throws Exception { + + if (provider.equals("Memory") || provider.equals("RocksDB")) { + // Local memory provider doesn't really have the concept of multiple sessions + return; + } + + @Cleanup + MetadataStoreExtended store1 = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().build()); + @Cleanup + MetadataStoreExtended store2 = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreConfig.builder().build()); + + @Cleanup + CoordinationService cs1 = new CoordinationServiceImpl(store1); + @Cleanup + LockManager lm1 = cs1.getLockManager(String.class); + + @Cleanup + CoordinationService cs2 = new CoordinationServiceImpl(store2); + @Cleanup + LockManager lm2 = cs2.getLockManager(String.class); + + String path1 = newKey(); + + ResourceLock lock1 = lm1.acquireLock(path1, "value-1").join(); + AtomicReference> lock2 = new AtomicReference<>(); + // lock 2 will steal the distributed lock first. + Awaitility.await().until(()-> { + // Ensure steal the lock success. + try { + lock2.set(lm2.acquireLock(path1, "value-1").join()); + return true; + } catch (Exception ex) { + return false; + } + }); + + // Since we can steal the lock repeatedly, we don't know which one will get it. + // But we can verify the final state. + Awaitility.await().untilAsserted(() -> { + if (lock1.getLockExpiredFuture().isDone()) { + assertTrue(lm1.listLocks(path1).join().isEmpty()); + assertFalse(lock2.get().getLockExpiredFuture().isDone()); + } else if (lock2.get().getLockExpiredFuture().isDone()) { + assertTrue(lm2.listLocks(path1).join().isEmpty()); + assertFalse(lock1.getLockExpiredFuture().isDone()); + } else { + fail("unexpected behaviour"); + } + }); + } } From a19e861af07e17862d3c77c96a2c55c3d6526910 Mon Sep 17 00:00:00 2001 From: lordcheng10 Date: Thu, 5 Jan 2023 12:52:25 +0800 Subject: [PATCH 3/4] update testCleanUpStateWhenRevalidationGotLockBusy --- .../java/org/apache/pulsar/metadata/LockManagerTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java index 0df5cadc679d6..0268d07c0ff82 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java @@ -232,7 +232,7 @@ public void revalidateLockOnDifferentSession(String provider, String url) throws } @Test(dataProvider = "impl") - public void testCleanUpStateWhenRevalidationGotLockBusy(String provider, Supplier urlSupplier) + public void testCleanUpStateWhenRevalidationGotLockBusy(String provider, String url) throws Exception { if (provider.equals("Memory") || provider.equals("RocksDB")) { @@ -241,10 +241,10 @@ public void testCleanUpStateWhenRevalidationGotLockBusy(String provider, Supplie } @Cleanup - MetadataStoreExtended store1 = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreExtended store1 = MetadataStoreExtended.create(url, MetadataStoreConfig.builder().build()); @Cleanup - MetadataStoreExtended store2 = MetadataStoreExtended.create(urlSupplier.get(), + MetadataStoreExtended store2 = MetadataStoreExtended.create(url, MetadataStoreConfig.builder().build()); @Cleanup From 5c583ecd53df33065b9223d51df67a829a42d808 Mon Sep 17 00:00:00 2001 From: lordcheng10 Date: Thu, 5 Jan 2023 15:00:49 +0800 Subject: [PATCH 4/4] update test testCleanUpStateWhenRevalidationGotLockBusy --- .../java/org/apache/pulsar/metadata/LockManagerTest.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java index 0268d07c0ff82..a7d0ed4454b86 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/LockManagerTest.java @@ -241,19 +241,16 @@ public void testCleanUpStateWhenRevalidationGotLockBusy(String provider, String } @Cleanup - MetadataStoreExtended store1 = MetadataStoreExtended.create(url, - MetadataStoreConfig.builder().build()); - @Cleanup - MetadataStoreExtended store2 = MetadataStoreExtended.create(url, + MetadataStoreExtended store = MetadataStoreExtended.create(url, MetadataStoreConfig.builder().build()); @Cleanup - CoordinationService cs1 = new CoordinationServiceImpl(store1); + CoordinationService cs1 = new CoordinationServiceImpl(store); @Cleanup LockManager lm1 = cs1.getLockManager(String.class); @Cleanup - CoordinationService cs2 = new CoordinationServiceImpl(store2); + CoordinationService cs2 = new CoordinationServiceImpl(store); @Cleanup LockManager lm2 = cs2.getLockManager(String.class);