diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java index e8230e0113ffe..7cbf594d58c00 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java @@ -227,7 +227,7 @@ public final CompletableFuture delete(String path, Optional expected } // Ensure caches are invalidated before the operation is confirmed return storeDelete(path, expectedVersion) - .thenRun(() -> { + .thenRunAsync(() -> { existsCache.synchronous().invalidate(path); String parent = parent(path); if (parent != null) { @@ -235,7 +235,7 @@ public final CompletableFuture delete(String path, Optional expected } metadataCaches.forEach(c -> c.invalidate(path)); - }); + }, executor); } @Override @@ -266,7 +266,7 @@ public final CompletableFuture put(String path, byte[] data, Optional { + .thenApplyAsync(stat -> { NotificationType type = stat.getVersion() == 0 ? NotificationType.Created : NotificationType.Modified; if (type == NotificationType.Created) { @@ -279,7 +279,7 @@ public final CompletableFuture put(String path, byte[] data, Optional c.refresh(path)); return stat; - }); + }, executor); } @Override diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/LocalMemoryMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/LocalMemoryMetadataStore.java index b80ebdbed622e..e4aeca97eb723 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/LocalMemoryMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/LocalMemoryMetadataStore.java @@ -23,6 +23,7 @@ import java.net.URISyntaxException; import java.util.ArrayList; import java.util.EnumSet; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.NavigableMap; @@ -64,6 +65,9 @@ private static class Value { private static final Map> STATIC_MAPS = new MapMaker() .weakValues().makeMap(); + // Manage all instances to facilitate registration to the same listener + private static final Map> STATIC_INSTANCE = new MapMaker() + .weakValues().makeMap(); private static final Map STATIC_ID_GEN_MAP = new MapMaker() .weakValues().makeMap(); @@ -84,6 +88,14 @@ public LocalMemoryMetadataStore(String metadataURL, MetadataStoreConfig metadata // Use a reference from a shared data set String name = uri.getHost(); map = STATIC_MAPS.computeIfAbsent(name, __ -> new TreeMap<>()); + STATIC_INSTANCE.compute(name, (key, value) -> { + if (value == null) { + value = new HashSet<>(); + } + value.forEach(v -> registerListener(v)); + value.add(this); + return value; + }); sequentialIdGenerator = STATIC_ID_GEN_MAP.computeIfAbsent(name, __ -> new AtomicLong()); log.info("Created LocalMemoryDataStore for '{}'", name); }