From 957cbe4b24719a4d4c72138d2c92ea5f88e9b154 Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Mon, 27 Jan 2025 12:07:22 -0800 Subject: [PATCH 01/14] [fix][metadata] fixed ephemeral zk put --- .../pulsar/metadata/impl/ZKMetadataStore.java | 6 ++-- .../metadata/MetadataStoreExtendedTest.java | 35 +++++++++++++++++++ 2 files changed, 38 insertions(+), 3 deletions(-) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java index 4c24aa5938b93..f0971791ae0ef 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/ZKMetadataStore.java @@ -439,8 +439,8 @@ private void internalStorePut(OpPut opPut) { future.completeExceptionally(getException(Code.BADVERSION, opPut.getPath())); } else { // The z-node does not exist, let's create it first - put(opPut.getPath(), opPut.getData(), Optional.of(-1L)).thenAccept( - s -> future.complete(s)) + put(opPut.getPath(), opPut.getData(), Optional.of(-1L), opPut.getOptions()) + .thenAccept(s -> future.complete(s)) .exceptionally(ex -> { if (ex.getCause() instanceof BadVersionException) { // The z-node exist now, let's overwrite it @@ -478,7 +478,7 @@ public void close() throws Exception { private Stat getStat(String path, org.apache.zookeeper.data.Stat zkStat) { return new Stat(path, zkStat.getVersion(), zkStat.getCtime(), zkStat.getMtime(), - zkStat.getEphemeralOwner() != -1, + zkStat.getEphemeralOwner() != 0, zkStat.getEphemeralOwner() == zkc.getSessionId()); } diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java index 9a38cdbcd2f85..25d7445685890 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.metadata; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; @@ -66,4 +67,38 @@ public void sequentialKeys(String provider, Supplier urlSupplier) throws assertNotEquals(seq1, seq2); assertTrue(n1 < n2); } + + @Test(dataProvider = "impl") + public void testPersistentOrEphemeralPut(String provider, Supplier urlSupplier) throws Exception { + final String key1 = newKey(); + MetadataStoreExtended store = MetadataStoreExtended.create(urlSupplier.get(), MetadataStoreConfig.builder().build()); + store.put(key1, "value-1".getBytes(), Optional.empty(), EnumSet.noneOf(CreateOption.class)).join(); + var value = store.get(key1).join().get(); + assertEquals(value.getValue(), "value-1".getBytes()); + assertFalse(value.getStat().isEphemeral()); + assertTrue(value.getStat().isFirstVersion()); + var version = value.getStat().getVersion(); + + store.put(key1, "value-2".getBytes(), Optional.empty(), EnumSet.noneOf(CreateOption.class)).join(); + value = store.get(key1).join().get(); + assertEquals(value.getValue(), "value-2".getBytes()); + assertFalse(value.getStat().isEphemeral()); + assertEquals(value.getStat().getVersion(), version + 1); + + final String key2 = newKey(); + store.put(key2, "value-4".getBytes(), Optional.empty(), EnumSet.of(CreateOption.Ephemeral)).join(); + value = store.get(key2).join().get(); + assertEquals(value.getValue(), "value-4".getBytes()); + assertTrue(value.getStat().isEphemeral()); + assertTrue(value.getStat().isFirstVersion()); + version = value.getStat().getVersion(); + + + store.put(key2, "value-5".getBytes(), Optional.empty(), EnumSet.of(CreateOption.Ephemeral)).join(); + value = store.get(key2).join().get(); + assertEquals(value.getValue(), "value-5".getBytes()); + assertTrue(value.getStat().isEphemeral()); + assertEquals(value.getStat().getVersion(), version + 1); + } + } From 2ca60952021f16eccb0af25ef61460441b55aed2 Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Mon, 27 Jan 2025 12:53:47 -0800 Subject: [PATCH 02/14] fixed format --- .../org/apache/pulsar/metadata/MetadataStoreExtendedTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java index 25d7445685890..b71511aabceae 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataStoreExtendedTest.java @@ -18,8 +18,8 @@ */ package org.apache.pulsar.metadata; -import static org.junit.jupiter.api.Assertions.assertFalse; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; From 4bd7516f3450e67a503ca63efb84e94612f583b7 Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Mon, 27 Jan 2025 17:43:44 -0800 Subject: [PATCH 03/14] fixed persistent lock --- .../metadata/coordination/impl/ResourceLockImpl.java | 11 +++++++++++ 1 file changed, 11 insertions(+) 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 692f224594cae..1015e6ad65a35 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 @@ -188,6 +188,17 @@ private CompletableFuture acquireWithNoRevalidation(T newValue) { CompletableFuture result = new CompletableFuture<>(); store.put(path, payload, Optional.of(version), EnumSet.of(CreateOption.Ephemeral)) .thenAccept(stat -> { + if (!stat.isEphemeral()) { + log.warn("Found persistent lock at {}. Trying to delete it and acquire it", path); + store.delete(path, Optional.of(stat.getVersion())) + .thenRun(() -> + // Reset the expectation that the key is not there anymore + setVersion(-1L) + ) + .thenCompose(__ -> acquireWithNoRevalidation(newValue)) + .thenRun(() -> log.info("Successfully recovered persistent lock at {}", path)); + return; + } synchronized (ResourceLockImpl.this) { state = State.Valid; version = stat.getVersion(); From 99f2d9d5da1fa81b22739b205add923e9e1fa543 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 30 Jan 2025 00:13:51 +0200 Subject: [PATCH 04/14] Use 0 as ephemeralOwner value of persistent nodes in MockZooKeeper --- .../main/java/org/apache/zookeeper/MockZooKeeper.java | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java index f32036e53f001..088a3be9ba8a2 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java @@ -64,6 +64,9 @@ import org.slf4j.LoggerFactory; public class MockZooKeeper extends ZooKeeper { + // ephemeralOwner value for persistent nodes + private static final long NOT_EPHEMERAL = 0L; + @Data @AllArgsConstructor private static class MockZNode { @@ -278,7 +281,7 @@ public String create(String path, byte[] data, List acl, CreateMode createM MockZNode.of(parentNode.getContent(), parentVersion + 1, parentNode.getEphemeralOwner())); } - tree.put(path, MockZNode.of(data, 0, createMode.isEphemeral() ? getEphemeralOwner() : -1L)); + tree.put(path, MockZNode.of(data, 0, createMode.isEphemeral() ? getEphemeralOwner() : NOT_EPHEMERAL)); toNotifyCreate.addAll(watchers.get(path)); @@ -373,7 +376,7 @@ public void create(final String path, final byte[] data, final List acl, Cr cb.processResult(KeeperException.Code.NONODE.intValue(), path, ctx, null); } else { tree.put(name, MockZNode.of(data, 0, - createMode != null && createMode.isEphemeral() ? getEphemeralOwner() : -1L)); + createMode != null && createMode.isEphemeral() ? getEphemeralOwner() : NOT_EPHEMERAL)); watchers.removeAll(name); unlockIfLocked(); cb.processResult(0, path, ctx, name); @@ -678,9 +681,7 @@ private static Stat createStatForZNode(MockZNode zNode) { private static Stat applyToStat(MockZNode zNode, Stat stat) { stat.setVersion(zNode.getVersion()); - if (zNode.getEphemeralOwner() != -1L) { - stat.setEphemeralOwner(zNode.getEphemeralOwner()); - } + stat.setEphemeralOwner(zNode.getEphemeralOwner()); return stat; } From 90aeec02c4ef9a526484b82c690c10474572630f Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 30 Jan 2025 00:23:47 +0200 Subject: [PATCH 05/14] Delete ephemeral znodes when a MockZooKeeperSession is closed --- .../main/java/org/apache/zookeeper/MockZooKeeper.java | 11 +++++++++++ .../org/apache/zookeeper/MockZooKeeperSession.java | 4 ++++ 2 files changed, 15 insertions(+) diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java index 088a3be9ba8a2..d64445b380f33 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java @@ -1200,5 +1200,16 @@ private void triggerPersistentWatches(String path, String parent, EventType even }); } + public void deleteEmpheralNodes(long sessionId) { + if (sessionId != NOT_EPHEMERAL) { + lock(); + try { + tree.values().removeIf(zNode -> zNode.getEphemeralOwner() == sessionId); + } finally { + unlockIfLocked(); + } + } + } + private static final Logger log = LoggerFactory.getLogger(MockZooKeeper.class); } diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java index a286a75aa9103..0e71a5f09de64 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java @@ -221,12 +221,16 @@ public void addWatch(String basePath, AddWatchMode mode, VoidCallback cb, Object public void close() throws InterruptedException { if (closeMockZooKeeperOnClose) { mockZooKeeper.close(); + } else { + mockZooKeeper.deleteEmpheralNodes(getSessionId()); } } public void shutdown() throws InterruptedException { if (closeMockZooKeeperOnClose) { mockZooKeeper.shutdown(); + } else { + mockZooKeeper.deleteEmpheralNodes(getSessionId()); } } From 9275d02ffc10e1f1e74d09c2d0c4ad8961585ba8 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 30 Jan 2025 00:33:10 +0200 Subject: [PATCH 06/14] Support passing ephemeral owner for multi ops in MockZooKeeperSession --- .../org/apache/zookeeper/MockZooKeeperSession.java | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java index 0e71a5f09de64..c890d48cf6b0c 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java @@ -188,12 +188,22 @@ public void delete(final String path, int version, final VoidCallback cb, final @Override public void multi(Iterable ops, AsyncCallback.MultiCallback cb, Object ctx) { - mockZooKeeper.multi(ops, cb, ctx); + try { + mockZooKeeper.overrideEpheralOwner(getSessionId()); + mockZooKeeper.multi(ops, cb, ctx); + } finally { + mockZooKeeper.removeEpheralOwnerOverride(); + } } @Override public List multi(Iterable ops) throws InterruptedException, KeeperException { - return mockZooKeeper.multi(ops); + try { + mockZooKeeper.overrideEpheralOwner(getSessionId()); + return mockZooKeeper.multi(ops); + } finally { + mockZooKeeper.removeEpheralOwnerOverride(); + } } @Override From 3e34b9235d355571a1fa77dc8ca05af260ac3024 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 30 Jan 2025 00:52:29 +0200 Subject: [PATCH 07/14] Fix typo --- .../src/main/java/org/apache/zookeeper/MockZooKeeper.java | 2 +- .../main/java/org/apache/zookeeper/MockZooKeeperSession.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java index d64445b380f33..6640bffa71268 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java @@ -1200,7 +1200,7 @@ private void triggerPersistentWatches(String path, String parent, EventType even }); } - public void deleteEmpheralNodes(long sessionId) { + public void deleteEphemeralNodes(long sessionId) { if (sessionId != NOT_EPHEMERAL) { lock(); try { diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java index c890d48cf6b0c..3b7224b2063f8 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java @@ -232,7 +232,7 @@ public void close() throws InterruptedException { if (closeMockZooKeeperOnClose) { mockZooKeeper.close(); } else { - mockZooKeeper.deleteEmpheralNodes(getSessionId()); + mockZooKeeper.deleteEphemeralNodes(getSessionId()); } } @@ -240,7 +240,7 @@ public void shutdown() throws InterruptedException { if (closeMockZooKeeperOnClose) { mockZooKeeper.shutdown(); } else { - mockZooKeeper.deleteEmpheralNodes(getSessionId()); + mockZooKeeper.deleteEphemeralNodes(getSessionId()); } } From e979320212a7f9fe6fd3bc648902be40ff86a53f Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 30 Jan 2025 00:55:07 +0200 Subject: [PATCH 08/14] Fix existing typos --- .../org/apache/zookeeper/MockZooKeeper.java | 18 +++++++++--------- .../apache/zookeeper/MockZooKeeperSession.java | 16 ++++++++-------- 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java index 6640bffa71268..a41c47e650e06 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java @@ -93,7 +93,7 @@ static MockZNode of(byte[] content, int version, long ephemeralOwner) { private ReentrantLock mutex; private AtomicLong sequentialIdGenerator; - private ThreadLocal epheralOwnerThreadLocal; + private ThreadLocal ephemeralOwnerThreadLocal; //see details of Objenesis caching - http://objenesis.org/details.html //see supported jvms - https://github.com/easymock/objenesis/blob/master/SupportedJVMs.md @@ -159,7 +159,7 @@ private static MockZooKeeper createMockZooKeeperInstance(ExecutorService executo ObjectInstantiator mockZooKeeperInstantiator = objenesis.getInstantiatorOf(MockZooKeeper.class); MockZooKeeper zk = mockZooKeeperInstantiator.newInstance(); - zk.epheralOwnerThreadLocal = new ThreadLocal<>(); + zk.ephemeralOwnerThreadLocal = new ThreadLocal<>(); zk.init(executor); zk.readOpDelayMs = readOpDelayMs; zk.mutex = new ReentrantLock(); @@ -313,19 +313,19 @@ public String create(String path, byte[] data, List acl, CreateMode createM } protected long getEphemeralOwner() { - Long epheralOwner = epheralOwnerThreadLocal.get(); - if (epheralOwner != null) { - return epheralOwner; + Long ephemeralOwner = ephemeralOwnerThreadLocal.get(); + if (ephemeralOwner != null) { + return ephemeralOwner; } return getSessionId(); } - public void overrideEpheralOwner(long epheralOwner) { - epheralOwnerThreadLocal.set(epheralOwner); + public void overrideEphemeralOwner(long ephemeralOwner) { + ephemeralOwnerThreadLocal.set(ephemeralOwner); } - public void removeEpheralOwnerOverride() { - epheralOwnerThreadLocal.remove(); + public void removeEphemeralOwnerOverride() { + ephemeralOwnerThreadLocal.remove(); } @Override diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java index 3b7224b2063f8..a75402018ac04 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeperSession.java @@ -88,10 +88,10 @@ public void register(Watcher watcher) { public String create(String path, byte[] data, List acl, CreateMode createMode) throws KeeperException, InterruptedException { try { - mockZooKeeper.overrideEpheralOwner(getSessionId()); + mockZooKeeper.overrideEphemeralOwner(getSessionId()); return mockZooKeeper.create(path, data, acl, createMode); } finally { - mockZooKeeper.removeEpheralOwnerOverride(); + mockZooKeeper.removeEphemeralOwnerOverride(); } } @@ -99,10 +99,10 @@ public String create(String path, byte[] data, List acl, CreateMode createM public void create(final String path, final byte[] data, final List acl, CreateMode createMode, final AsyncCallback.StringCallback cb, final Object ctx) { try { - mockZooKeeper.overrideEpheralOwner(getSessionId()); + mockZooKeeper.overrideEphemeralOwner(getSessionId()); mockZooKeeper.create(path, data, acl, createMode, cb, ctx); } finally { - mockZooKeeper.removeEpheralOwnerOverride(); + mockZooKeeper.removeEphemeralOwnerOverride(); } } @@ -189,20 +189,20 @@ public void delete(final String path, int version, final VoidCallback cb, final @Override public void multi(Iterable ops, AsyncCallback.MultiCallback cb, Object ctx) { try { - mockZooKeeper.overrideEpheralOwner(getSessionId()); + mockZooKeeper.overrideEphemeralOwner(getSessionId()); mockZooKeeper.multi(ops, cb, ctx); } finally { - mockZooKeeper.removeEpheralOwnerOverride(); + mockZooKeeper.removeEphemeralOwnerOverride(); } } @Override public List multi(Iterable ops) throws InterruptedException, KeeperException { try { - mockZooKeeper.overrideEpheralOwner(getSessionId()); + mockZooKeeper.overrideEphemeralOwner(getSessionId()); return mockZooKeeper.multi(ops); } finally { - mockZooKeeper.removeEpheralOwnerOverride(); + mockZooKeeper.removeEphemeralOwnerOverride(); } } From 0e5a431b7d827f865d0a90dfda463810af19e955 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 30 Jan 2025 00:58:55 +0200 Subject: [PATCH 09/14] Make default session id 1L in MockZooKeeper --- testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java index a41c47e650e06..7b23d0814aab9 100644 --- a/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java +++ b/testmocks/src/main/java/org/apache/zookeeper/MockZooKeeper.java @@ -87,7 +87,7 @@ static MockZNode of(byte[] content, int version, long ephemeralOwner) { private ExecutorService executor; private Watcher sessionWatcher; - private long sessionId = 0L; + private long sessionId = 1L; private int readOpDelayMs; private ReentrantLock mutex; From b76c7cd50934703d507b6dd96e57b1c5230d1976 Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Sun, 9 Feb 2025 18:15:35 -0800 Subject: [PATCH 10/14] removed persistent lock check --- .../util/EntryWrapperPerformanceTestNG.java | 256 ++++++++++++++++++ .../coordination/impl/ResourceLockImpl.java | 11 - 2 files changed, 256 insertions(+), 11 deletions(-) create mode 100644 managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java new file mode 100644 index 0000000000000..ba661e0119654 --- /dev/null +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java @@ -0,0 +1,256 @@ +package org.apache.bookkeeper.mledger.util; + +import io.netty.util.Recycler; +import java.util.Collections; +import java.util.UUID; +import java.util.concurrent.locks.StampedLock; +import org.apache.commons.lang.RandomStringUtils; +import org.testng.annotations.Test; +import java.lang.management.GarbageCollectorMXBean; +import java.lang.management.ManagementFactory; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.*; + + +public class EntryWrapperPerformanceTestNG { + + private static class EntryWrapper { + private final Recycler.Handle recyclerHandle; + private static final Recycler RECYCLER = new Recycler() { + @Override + protected EntryWrapper newObject(Handle recyclerHandle) { + return new EntryWrapper(recyclerHandle); + } + }; + private final StampedLock lock = new StampedLock(); + private K key; + private V value; + long size; + + private EntryWrapper(Recycler.Handle recyclerHandle) { + this.recyclerHandle = recyclerHandle; + } + + static EntryWrapper create(K key, V value, long size) { + EntryWrapper entryWrapper = RECYCLER.get(); + long stamp = entryWrapper.lock.writeLock(); + entryWrapper.key = key; + entryWrapper.value = value; + entryWrapper.size = size; + entryWrapper.lock.unlockWrite(stamp); + return entryWrapper; + } + + K getKey() { + long stamp = lock.tryOptimisticRead(); + K localKey = key; + if (!lock.validate(stamp)) { + stamp = lock.readLock(); + localKey = key; + lock.unlockRead(stamp); + } + return localKey; + } + + V getValue(K key) { + long stamp = lock.tryOptimisticRead(); + K localKey = this.key; + V localValue = this.value; + if (!lock.validate(stamp)) { + stamp = lock.readLock(); + localKey = this.key; + localValue = this.value; + lock.unlockRead(stamp); + } + if (localKey != key) { + return null; + } + return localValue; + } + + long getSize() { + long stamp = lock.tryOptimisticRead(); + long localSize = size; + if (!lock.validate(stamp)) { + stamp = lock.readLock(); + localSize = size; + lock.unlockRead(stamp); + } + return localSize; + } + + void recycle() { + key = null; + value = null; + size = 0; + recyclerHandle.recycle(this); + } + } + + + private static final int ITERATIONS = 1_000_000; + private static final int THREAD_COUNT = 10; // Simulating 10 concurrent threads + private static final int TEST_RUNS = 100; // Run each test 10 times + + private long getGCCount() { + long count = 0; + for (GarbageCollectorMXBean gcBean : ManagementFactory.getGarbageCollectorMXBeans()) { + long gcCount = gcBean.getCollectionCount(); + if (gcCount != -1) { + count += gcCount; + } + } + return count; + } + + private long getUsedMemory() { + System.gc(); // Suggest GC before measuring memory usage + return Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory(); + } + + private void runTest(String testName, Runnable testLogic) { + List times = new ArrayList<>(); + List gcCounts = new ArrayList<>(); + List memoryUsages = new ArrayList<>(); + + for (int i = 0; i < TEST_RUNS; i++) { + long startGC = getGCCount(); + long startMemory = getUsedMemory(); + long startTime = System.nanoTime(); + + testLogic.run(); + + long elapsedTime = System.nanoTime() - startTime; + long endGC = getGCCount(); + long endMemory = getUsedMemory(); + + times.add(elapsedTime); + gcCounts.add(endGC - startGC); + memoryUsages.add(endMemory - startMemory); + + System.out.println(testName + " - Run " + (i + 1) + " - Time: " + TimeUnit.NANOSECONDS.toMillis(elapsedTime) + " ms, GC: " + (endGC - startGC) + ", Memory: " + (endMemory - startMemory) + " bytes"); + } + + printStatistics(testName, times, gcCounts, memoryUsages); + } + + private void printStatistics(String testName, List times, List gcCounts, List memoryUsages) { + System.out.println("========= " + testName + " Summary ========="); + System.out.println("Average Time: " + TimeUnit.NANOSECONDS.toMillis(average(times)) + " ms"); + System.out.println("Median Time: " + TimeUnit.NANOSECONDS.toMillis(median(times)) + " ms"); + System.out.println("P99 Time: " + TimeUnit.NANOSECONDS.toMillis(percentile(times, 99)) + " ms"); + System.out.println("Max Time: " + TimeUnit.NANOSECONDS.toMillis(Collections.max(times)) + " ms"); + + System.out.println("Average GC: " + average(gcCounts)); + System.out.println("Median GC: " + median(gcCounts)); + System.out.println("P99 GC: " + percentile(gcCounts, 99)); + System.out.println("Max GC: " + Collections.max(gcCounts)); + + System.out.println("Average Memory: " + average(memoryUsages) + " bytes"); + System.out.println("Median Memory: " + median(memoryUsages) + " bytes"); + System.out.println("P99 Memory: " + percentile(memoryUsages, 99) + " bytes"); + System.out.println("Max Memory: " + Collections.max(memoryUsages) + " bytes"); + System.out.println("===================================="); + } + + private long average(List values) { + return values.stream().mapToLong(Long::longValue).sum() / values.size(); + } + + private long median(List values) { + List sorted = new ArrayList<>(values); + Collections.sort(sorted); + int middle = sorted.size() / 2; + return (sorted.size() % 2 == 0) ? (sorted.get(middle - 1) + sorted.get(middle)) / 2 : sorted.get(middle); + } + + private long percentile(List values, int percentile) { + List sorted = new ArrayList<>(values); + Collections.sort(sorted); + int index = (int) Math.ceil(percentile / 100.0 * sorted.size()) - 1; + return sorted.get(Math.max(0, Math.min(index, sorted.size() - 1))); + } + + private static final String CHARACTERS = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789"; + + public static String generateString() { + StringBuilder sb = new StringBuilder(5); + for (int i = 0; i < 5; i++) { + int index = ThreadLocalRandom.current().nextInt(CHARACTERS.length()); + sb.append(CHARACTERS.charAt(index)); + } + return sb.toString(); + } + + @Test + public void testObjectPoolingPerformance() { + runTest("Object Pooling", () -> { + for (int i = 0; i < ITERATIONS; i++) { + EntryWrapper entry = EntryWrapper.create(i, generateString(), i); + entry.recycle(); + } + }); + } + + @Test + public void testNewObjectPerformance() { + runTest("New Object Creation", () -> { + for (int i = 0; i < ITERATIONS; i++) { + var str = generateString(); // Creates a new String object each time + } + }); + } + + @Test + public void testConcurrentObjectPooling() { + runTest("Concurrent Pooling", () -> { + ExecutorService executor = Executors.newFixedThreadPool(THREAD_COUNT); + List> futures = new ArrayList<>(); + + for (int t = 0; t < THREAD_COUNT; t++) { + futures.add(executor.submit(() -> { + for (int i = 0; i < ITERATIONS / THREAD_COUNT; i++) { + EntryWrapper entry = EntryWrapper.create(i, generateString(), i); + entry.recycle(); + } + })); + } + + futures.forEach(f -> { + try { + f.get(); + } catch (Exception ignored) {} + }); + + executor.shutdown(); + }); + } + + @Test + public void testConcurrentNewObjectCreation() { + runTest("Concurrent New Object", () -> { + ExecutorService executor = Executors.newFixedThreadPool(THREAD_COUNT); + List> futures = new ArrayList<>(); + + for (int t = 0; t < THREAD_COUNT; t++) { + futures.add(executor.submit(() -> { + //List strings = new ArrayList<>(); + for (int i = 0; i < ITERATIONS / THREAD_COUNT; i++) { + var str = generateString(); // Creates a new String object each time + } + //strings.clear(); + })); + } + + futures.forEach(f -> { + try { + f.get(); + } catch (Exception ignored) {} + }); + + executor.shutdown(); + }); + } +} + 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 1015e6ad65a35..692f224594cae 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 @@ -188,17 +188,6 @@ private CompletableFuture acquireWithNoRevalidation(T newValue) { CompletableFuture result = new CompletableFuture<>(); store.put(path, payload, Optional.of(version), EnumSet.of(CreateOption.Ephemeral)) .thenAccept(stat -> { - if (!stat.isEphemeral()) { - log.warn("Found persistent lock at {}. Trying to delete it and acquire it", path); - store.delete(path, Optional.of(stat.getVersion())) - .thenRun(() -> - // Reset the expectation that the key is not there anymore - setVersion(-1L) - ) - .thenCompose(__ -> acquireWithNoRevalidation(newValue)) - .thenRun(() -> log.info("Successfully recovered persistent lock at {}", path)); - return; - } synchronized (ResourceLockImpl.this) { state = State.Valid; version = stat.getVersion(); From 313ccb7b95071be3232bc5e41b01147cb7777faf Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Sun, 9 Feb 2025 18:23:57 -0800 Subject: [PATCH 11/14] Removed unused file --- .../util/EntryWrapperPerformanceTestNG.java | 256 ------------------ 1 file changed, 256 deletions(-) delete mode 100644 managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java deleted file mode 100644 index ba661e0119654..0000000000000 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/EntryWrapperPerformanceTestNG.java +++ /dev/null @@ -1,256 +0,0 @@ -package org.apache.bookkeeper.mledger.util; - -import io.netty.util.Recycler; -import java.util.Collections; -import java.util.UUID; -import java.util.concurrent.locks.StampedLock; -import org.apache.commons.lang.RandomStringUtils; -import org.testng.annotations.Test; -import java.lang.management.GarbageCollectorMXBean; -import java.lang.management.ManagementFactory; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.*; - - -public class EntryWrapperPerformanceTestNG { - - private static class EntryWrapper { - private final Recycler.Handle recyclerHandle; - private static final Recycler RECYCLER = new Recycler() { - @Override - protected EntryWrapper newObject(Handle recyclerHandle) { - return new EntryWrapper(recyclerHandle); - } - }; - private final StampedLock lock = new StampedLock(); - private K key; - private V value; - long size; - - private EntryWrapper(Recycler.Handle recyclerHandle) { - this.recyclerHandle = recyclerHandle; - } - - static EntryWrapper create(K key, V value, long size) { - EntryWrapper entryWrapper = RECYCLER.get(); - long stamp = entryWrapper.lock.writeLock(); - entryWrapper.key = key; - entryWrapper.value = value; - entryWrapper.size = size; - entryWrapper.lock.unlockWrite(stamp); - return entryWrapper; - } - - K getKey() { - long stamp = lock.tryOptimisticRead(); - K localKey = key; - if (!lock.validate(stamp)) { - stamp = lock.readLock(); - localKey = key; - lock.unlockRead(stamp); - } - return localKey; - } - - V getValue(K key) { - long stamp = lock.tryOptimisticRead(); - K localKey = this.key; - V localValue = this.value; - if (!lock.validate(stamp)) { - stamp = lock.readLock(); - localKey = this.key; - localValue = this.value; - lock.unlockRead(stamp); - } - if (localKey != key) { - return null; - } - return localValue; - } - - long getSize() { - long stamp = lock.tryOptimisticRead(); - long localSize = size; - if (!lock.validate(stamp)) { - stamp = lock.readLock(); - localSize = size; - lock.unlockRead(stamp); - } - return localSize; - } - - void recycle() { - key = null; - value = null; - size = 0; - recyclerHandle.recycle(this); - } - } - - - private static final int ITERATIONS = 1_000_000; - private static final int THREAD_COUNT = 10; // Simulating 10 concurrent threads - private static final int TEST_RUNS = 100; // Run each test 10 times - - private long getGCCount() { - long count = 0; - for (GarbageCollectorMXBean gcBean : ManagementFactory.getGarbageCollectorMXBeans()) { - long gcCount = gcBean.getCollectionCount(); - if (gcCount != -1) { - count += gcCount; - } - } - return count; - } - - private long getUsedMemory() { - System.gc(); // Suggest GC before measuring memory usage - return Runtime.getRuntime().totalMemory() - Runtime.getRuntime().freeMemory(); - } - - private void runTest(String testName, Runnable testLogic) { - List times = new ArrayList<>(); - List gcCounts = new ArrayList<>(); - List memoryUsages = new ArrayList<>(); - - for (int i = 0; i < TEST_RUNS; i++) { - long startGC = getGCCount(); - long startMemory = getUsedMemory(); - long startTime = System.nanoTime(); - - testLogic.run(); - - long elapsedTime = System.nanoTime() - startTime; - long endGC = getGCCount(); - long endMemory = getUsedMemory(); - - times.add(elapsedTime); - gcCounts.add(endGC - startGC); - memoryUsages.add(endMemory - startMemory); - - System.out.println(testName + " - Run " + (i + 1) + " - Time: " + TimeUnit.NANOSECONDS.toMillis(elapsedTime) + " ms, GC: " + (endGC - startGC) + ", Memory: " + (endMemory - startMemory) + " bytes"); - } - - printStatistics(testName, times, gcCounts, memoryUsages); - } - - private void printStatistics(String testName, List times, List gcCounts, List memoryUsages) { - System.out.println("========= " + testName + " Summary ========="); - System.out.println("Average Time: " + TimeUnit.NANOSECONDS.toMillis(average(times)) + " ms"); - System.out.println("Median Time: " + TimeUnit.NANOSECONDS.toMillis(median(times)) + " ms"); - System.out.println("P99 Time: " + TimeUnit.NANOSECONDS.toMillis(percentile(times, 99)) + " ms"); - System.out.println("Max Time: " + TimeUnit.NANOSECONDS.toMillis(Collections.max(times)) + " ms"); - - System.out.println("Average GC: " + average(gcCounts)); - System.out.println("Median GC: " + median(gcCounts)); - System.out.println("P99 GC: " + percentile(gcCounts, 99)); - System.out.println("Max GC: " + Collections.max(gcCounts)); - - System.out.println("Average Memory: " + average(memoryUsages) + " bytes"); - System.out.println("Median Memory: " + median(memoryUsages) + " bytes"); - System.out.println("P99 Memory: " + percentile(memoryUsages, 99) + " bytes"); - System.out.println("Max Memory: " + Collections.max(memoryUsages) + " bytes"); - System.out.println("===================================="); - } - - private long average(List values) { - return values.stream().mapToLong(Long::longValue).sum() / values.size(); - } - - private long median(List values) { - List sorted = new ArrayList<>(values); - Collections.sort(sorted); - int middle = sorted.size() / 2; - return (sorted.size() % 2 == 0) ? (sorted.get(middle - 1) + sorted.get(middle)) / 2 : sorted.get(middle); - } - - private long percentile(List values, int percentile) { - List sorted = new ArrayList<>(values); - Collections.sort(sorted); - int index = (int) Math.ceil(percentile / 100.0 * sorted.size()) - 1; - return sorted.get(Math.max(0, Math.min(index, sorted.size() - 1))); - } - - private static final String CHARACTERS = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789"; - - public static String generateString() { - StringBuilder sb = new StringBuilder(5); - for (int i = 0; i < 5; i++) { - int index = ThreadLocalRandom.current().nextInt(CHARACTERS.length()); - sb.append(CHARACTERS.charAt(index)); - } - return sb.toString(); - } - - @Test - public void testObjectPoolingPerformance() { - runTest("Object Pooling", () -> { - for (int i = 0; i < ITERATIONS; i++) { - EntryWrapper entry = EntryWrapper.create(i, generateString(), i); - entry.recycle(); - } - }); - } - - @Test - public void testNewObjectPerformance() { - runTest("New Object Creation", () -> { - for (int i = 0; i < ITERATIONS; i++) { - var str = generateString(); // Creates a new String object each time - } - }); - } - - @Test - public void testConcurrentObjectPooling() { - runTest("Concurrent Pooling", () -> { - ExecutorService executor = Executors.newFixedThreadPool(THREAD_COUNT); - List> futures = new ArrayList<>(); - - for (int t = 0; t < THREAD_COUNT; t++) { - futures.add(executor.submit(() -> { - for (int i = 0; i < ITERATIONS / THREAD_COUNT; i++) { - EntryWrapper entry = EntryWrapper.create(i, generateString(), i); - entry.recycle(); - } - })); - } - - futures.forEach(f -> { - try { - f.get(); - } catch (Exception ignored) {} - }); - - executor.shutdown(); - }); - } - - @Test - public void testConcurrentNewObjectCreation() { - runTest("Concurrent New Object", () -> { - ExecutorService executor = Executors.newFixedThreadPool(THREAD_COUNT); - List> futures = new ArrayList<>(); - - for (int t = 0; t < THREAD_COUNT; t++) { - futures.add(executor.submit(() -> { - //List strings = new ArrayList<>(); - for (int i = 0; i < ITERATIONS / THREAD_COUNT; i++) { - var str = generateString(); // Creates a new String object each time - } - //strings.clear(); - })); - } - - futures.forEach(f -> { - try { - f.get(); - } catch (Exception ignored) {} - }); - - executor.shutdown(); - }); - } -} - From 17ec37a13e32443c14bfa4647a2cf03f66b9aab0 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 12 Feb 2025 00:34:37 +0200 Subject: [PATCH 12/14] Add .withTestZookeeper to PulsarTestContext --- .../auth/MockedPulsarServiceBaseTest.java | 20 ++++- .../broker/testcontext/PulsarTestContext.java | 75 +++++++++++++++++-- 2 files changed, 89 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index 42e2c00f73acf..81d1a105c175f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -150,6 +150,9 @@ public static String getTlsFileForClient(String name) { private final List closeables = new ArrayList<>(); + // Set to true in test's constructor to use a real Zookeeper (TestZKServer) + protected boolean useTestZookeeper; + public MockedPulsarServiceBaseTest() { resetConfig(); } @@ -461,7 +464,6 @@ protected PulsarTestContext.Builder createPulsarTestContextBuilder(ServiceConfig PulsarTestContext.Builder builder = PulsarTestContext.builder() .spyByDefault() .config(conf) - .withMockZookeeper(true) .pulsarServiceCustomizer(pulsarService -> { try { beforePulsarStart(pulsarService); @@ -470,9 +472,25 @@ protected PulsarTestContext.Builder createPulsarTestContextBuilder(ServiceConfig } }) .brokerServiceCustomizer(this::customizeNewBrokerService); + configureMetadataStores(builder); return builder; } + /** + * Configures the metadata stores for the PulsarTestContext.Builder instance. + * Set useTestZookeeper to true in the test's constructor to use TestZKServer which is a real ZooKeeper + * implementation. + * + * @param builder the PulsarTestContext.Builder instance to configure + */ + protected void configureMetadataStores(PulsarTestContext.Builder builder) { + if (useTestZookeeper) { + builder.withTestZookeeper(); + } else { + builder.withMockZookeeper(); + } + } + protected PulsarTestContext createAdditionalPulsarTestContext(ServiceConfiguration conf) throws Exception { return createAdditionalPulsarTestContext(conf, null); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java index 6403c3bcec4c3..ef3c841ef1d24 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java @@ -65,6 +65,7 @@ import org.apache.pulsar.compaction.CompactionServiceFactory; import org.apache.pulsar.compaction.Compactor; import org.apache.pulsar.compaction.PulsarCompactionServiceFactory; +import org.apache.pulsar.metadata.TestZKServer; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.MetadataStoreConfig; import org.apache.pulsar.metadata.api.MetadataStoreException; @@ -72,8 +73,11 @@ import org.apache.pulsar.metadata.impl.MetadataStoreFactoryImpl; import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.KeeperException; import org.apache.zookeeper.MockZooKeeper; import org.apache.zookeeper.MockZooKeeperSession; +import org.apache.zookeeper.ZooDefs; +import org.apache.zookeeper.ZooKeeper; import org.apache.zookeeper.data.ACL; import org.jetbrains.annotations.NotNull; import org.mockito.Mockito; @@ -477,16 +481,77 @@ public Builder withMockZookeeper(boolean useSeparateGlobalZk) { private MockZooKeeper createMockZooKeeper() throws Exception { MockZooKeeper zk = MockZooKeeper.newInstance(MoreExecutors.newDirectExecutorService()); - List dummyAclList = new ArrayList<>(0); + initializeZookeeper(zk); + registerCloseable(zk::shutdown); + return zk; + } + private static void initializeZookeeper(ZooKeeper zk) throws KeeperException, InterruptedException { ZkUtils.createFullPathOptimistic(zk, "/ledgers/available/192.168.1.1:" + 5000, - "".getBytes(StandardCharsets.UTF_8), dummyAclList, CreateMode.PERSISTENT); + "".getBytes(StandardCharsets.UTF_8), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - zk.create("/ledgers/LAYOUT", "1\nflat:1".getBytes(StandardCharsets.UTF_8), dummyAclList, + zk.create("/ledgers/LAYOUT", "1\nflat:1".getBytes(StandardCharsets.UTF_8), ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + } - registerCloseable(zk::shutdown); - return zk; + /** + * Configure this PulsarTestContext to use a test ZooKeeper instance which is + * shared for both the local and configuration metadata stores. + * + * @return the builder + */ + public Builder withTestZookeeper() { + return withTestZookeeper(false); + } + + /** + * Configure this PulsarTestContext to use a test ZooKeeper instance. + * + * @param useSeparateGlobalZk if true, the global (configuration) zookeeper will be a separate instance + * @return the builder + */ + public Builder withTestZookeeper(boolean useSeparateGlobalZk) { + try { + TestZKServer localZk = createTestZookeeper(); + MetadataStoreExtended localStore = + createTestZookeeperMetadataStore(localZk, MetadataStoreConfig.METADATA_STORE); + localMetadataStore(localStore); + MetadataStoreExtended configStore; + if (useSeparateGlobalZk) { + TestZKServer globalZk = createTestZookeeper(); + configStore = createTestZookeeperMetadataStore(globalZk, + MetadataStoreConfig.CONFIGURATION_METADATA_STORE); + } else { + configStore = + createTestZookeeperMetadataStore(localZk, MetadataStoreConfig.CONFIGURATION_METADATA_STORE); + } + configurationMetadataStore(configStore); + } catch (Exception e) { + throw new RuntimeException(e); + } + return this; + } + + private TestZKServer createTestZookeeper() throws Exception { + TestZKServer testZKServer = new TestZKServer(); + try (ZooKeeper zkc = new ZooKeeper(testZKServer.getConnectionString(), 5000, event -> { + })) { + initializeZookeeper(zkc); + } + registerCloseable(testZKServer); + return testZKServer; + } + + private MetadataStoreExtended createTestZookeeperMetadataStore(TestZKServer zkServer, + String metadataStoreName) { + try { + MetadataStoreExtended store = MetadataStoreExtended.create("zk:" + zkServer.getConnectionString(), + MetadataStoreConfig.builder().metadataStoreName(metadataStoreName).build()); + registerCloseable(store); + return store; + } catch (MetadataStoreException e) { + throw new RuntimeException(e); + } } /** From 083d3edc93f5cf57038d0ec9ada9510073226354 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 12 Feb 2025 01:18:38 +0200 Subject: [PATCH 13/14] Run BrokerServiceLookupTest using real ZooKeeper and MockZooKeeper --- .../client/api/BrokerServiceLookupTest.java | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java index 07deb9007c487..34db053271e1a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/BrokerServiceLookupTest.java @@ -112,14 +112,28 @@ import org.awaitility.reflect.WhiteboxImpl; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.testng.SkipException; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Factory; import org.testng.annotations.Test; @Test(groups = "broker-api") public class BrokerServiceLookupTest extends ProducerConsumerBase { private static final Logger log = LoggerFactory.getLogger(BrokerServiceLookupTest.class); + @DataProvider + private static Object[] booleanValues() { + return new Object[]{ true, false }; + } + + @Factory(dataProvider = "booleanValues") + public BrokerServiceLookupTest(boolean useTestZookeeper) { + // when set to true, TestZKServer is used which is a real ZooKeeper implementation + this.useTestZookeeper = useTestZookeeper; + } + @BeforeMethod @Override protected void setup() throws Exception { @@ -1197,6 +1211,9 @@ public String authenticate(AuthenticationDataSource authData) throws Authenticat @Test public void testLookupConnectionNotCloseIfGetUnloadingExOrMetadataEx() throws Exception { + if (useTestZookeeper) { + throw new SkipException("This test case depends on MockZooKeeper"); + } String tpName = BrokerTestUtil.newUniqueName("persistent://public/default/tp"); admin.topics().createNonPartitionedTopic(tpName); PulsarClientImpl pulsarClientImpl = (PulsarClientImpl) pulsarClient; From ca6cf8a3cd331e5abf1f7ff4955fda018de4350c Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 12 Feb 2025 01:26:32 +0200 Subject: [PATCH 14/14] Fix checkstyle --- .../org/apache/pulsar/broker/testcontext/PulsarTestContext.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java index ef3c841ef1d24..82a1b574ac921 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/testcontext/PulsarTestContext.java @@ -27,7 +27,6 @@ import java.io.IOException; import java.nio.charset.StandardCharsets; import java.time.Duration; -import java.util.ArrayList; import java.util.Collection; import java.util.List; import java.util.Optional; @@ -78,7 +77,6 @@ import org.apache.zookeeper.MockZooKeeperSession; import org.apache.zookeeper.ZooDefs; import org.apache.zookeeper.ZooKeeper; -import org.apache.zookeeper.data.ACL; import org.jetbrains.annotations.NotNull; import org.mockito.Mockito; import org.mockito.internal.util.MockUtil;