From b49f3031de6b9df8f914e646f0f2005b5f2b68f8 Mon Sep 17 00:00:00 2001 From: maxlisongsong Date: Thu, 16 Jul 2026 23:01:41 +0800 Subject: [PATCH] [fix][meta] Complete handleMetadataEvent future exceptionally when the initial get fails --- .../metadata/impl/AbstractMetadataStore.java | 10 ++++-- .../impl/MetadataEventSynchronizerTest.java | 33 +++++++++++++++++++ 2 files changed, 40 insertions(+), 3 deletions(-) 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 a8e11fc618093..a56cd25010b94 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 @@ -235,14 +235,14 @@ protected long getChildrenCacheMaxSizeBytes() { @Override public CompletableFuture handleMetadataEvent(MetadataEvent event) { CompletableFuture result = new CompletableFuture<>(); - get(event.getPath()).thenApply(res -> { + get(event.getPath()).thenAccept(res -> { Set options = event.getOptions() != null ? event.getOptions() : Collections.emptySet(); if (res.isPresent()) { GetResult existingValue = res.get(); if (shouldIgnoreEvent(event, existingValue)) { result.complete(null); - return result; + return; } } // else update the event @@ -262,7 +262,11 @@ public CompletableFuture handleMetadataEvent(MetadataEvent event) { } return false; }); - return result; + }).exceptionally(ex -> { + Throwable cause = FutureUtil.unwrapCompletionException(ex); + log.warn().attr("path", event.getPath()).exception(cause).log("Failed to handle metadata event"); + result.completeExceptionally(cause); + return null; }); return result; } diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/MetadataEventSynchronizerTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/MetadataEventSynchronizerTest.java index 3b07e0b3c2bd6..35531314df03e 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/MetadataEventSynchronizerTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/impl/MetadataEventSynchronizerTest.java @@ -20,13 +20,24 @@ import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; +import static org.testng.Assert.expectThrows; import java.nio.charset.StandardCharsets; +import java.util.HashSet; import java.util.Optional; +import java.util.Set; import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; import lombok.Cleanup; +import org.apache.pulsar.metadata.api.GetResult; +import org.apache.pulsar.metadata.api.MetadataEvent; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.MetadataStoreConfig; +import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.MetadataStoreFactory; +import org.apache.pulsar.metadata.api.NotificationType; +import org.apache.pulsar.metadata.api.Option; import org.awaitility.Awaitility; import org.testng.annotations.Test; @@ -75,6 +86,28 @@ public void testSharedInstance() throws Exception { }); } + @Test + public void testHandleMetadataEventCompletesWhenGetFails() throws Exception { + @Cleanup + LocalMemoryMetadataStore store = new LocalMemoryMetadataStore("memory:local", + MetadataStoreConfig.builder().build()) { + @Override + public CompletableFuture> storeGet(String path, Set