diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java index 2cbe9a6dc19b3..f58530bde311b 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/cache/impl/MetadataCacheImpl.java @@ -294,12 +294,12 @@ public void accept(Notification t) { private CompletableFuture executeWithRetry(Supplier> op, String key) { CompletableFuture result = new CompletableFuture<>(); - op.get().thenAccept(r -> result.complete(r)).exceptionally((ex) -> { + op.get().thenAccept(result::complete).exceptionally((ex) -> { if (ex.getCause() instanceof BadVersionException) { // if resource is updated by other than metadata-cache then metadata-cache will get bad-version // exception. so, try to invalidate the cache and try one more time. objCache.synchronous().invalidate(key); - op.get().thenAccept((c) -> result.complete(null)).exceptionally((ex1) -> { + op.get().thenAccept(result::complete).exceptionally((ex1) -> { result.completeExceptionally(ex1.getCause()); return null; }); diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java index be6a03d0eac38..43af3ad757ee4 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/MetadataCacheTest.java @@ -491,15 +491,21 @@ public void readModifyUpdateBadVersionRetry() throws Exception { MyClass value1 = new MyClass("a", 1); objCache1.create(key1, value1).join(); - objCache1.get(key1).join(); + assertEquals(objCache1.get(key1).join().get().b, 1); - objCache2.readModifyUpdate(key1, v -> { + CompletableFuture future1 = objCache1.readModifyUpdate(key1, v -> { return new MyClass(v.a, v.b + 1); - }).join(); + }); - objCache1.readModifyUpdate(key1, v -> { + CompletableFuture future2 = objCache2.readModifyUpdate(key1, v -> { return new MyClass(v.a, v.b + 1); - }).join(); + }); + + MyClass myClass1 = future1.join(); + assertEquals(myClass1.b, 2); + + MyClass myClass2 = future2.join(); + assertEquals(myClass2.b, 3); } @Test(dataProvider = "impl")