This is an automated email from the ASF dual-hosted git repository.

lhotari pushed a commit to branch branch-3.0
in repository https://gitbox.apache.org/repos/asf/pulsar.git


The following commit(s) were added to refs/heads/branch-3.0 by this push:
     new e097ea14024 [fix][broker] Consumer stuck when delete subscription 
__compaction failed (#23980)
e097ea14024 is described below

commit e097ea14024a2217eb454180b0b8d7194f33c8ce
Author: fengyubiao <[email protected]>
AuthorDate: Thu Apr 10 22:45:56 2025 +0800

    [fix][broker] Consumer stuck when delete subscription __compaction failed 
(#23980)
    
    Co-authored-by: Lari Hotari <[email protected]>
    (cherry picked from commit 98c99830ccbe208fbee3aedf1e55901cd31f942f)
---
 .../broker/service/persistent/PersistentTopic.java |  49 +++++--
 .../pulsar/compaction/CompactedTopicImpl.java      |  11 +-
 .../persistent/CompactionConcurrencyTest.java      | 153 +++++++++++++++++++++
 .../compaction/GetLastMessageIdCompactedTest.java  |  76 +++++++++-
 4 files changed, 276 insertions(+), 13 deletions(-)

diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
index 4e71237e30c..205eb98f5d6 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
@@ -229,9 +229,10 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
     protected final MessageDeduplication messageDeduplication;
 
     private static final long COMPACTION_NEVER_RUN = -0xfebecffeL;
-    private volatile CompletableFuture<Long> currentCompaction = 
CompletableFuture.completedFuture(
+    volatile CompletableFuture<Long> currentCompaction = 
CompletableFuture.completedFuture(
             COMPACTION_NEVER_RUN);
     private final CompactedTopic compactedTopic;
+    final AtomicBoolean disablingCompaction = new AtomicBoolean(false);
 
     // TODO: Create compaction strategy from topic policy when exposing 
strategic compaction to users.
     private static Map<String, TopicCompactionStrategy> strategicCompactionMap 
= Map.of(
@@ -1305,18 +1306,42 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
             return;
         }
 
-        currentCompaction.handle((__, e) -> {
-            if (e != null) {
-                log.warn("[{}][{}] Last compaction task failed", topic, 
subscriptionName);
+        // Avoid concurrently execute compaction and unsubscribing.
+        synchronized (this) {
+            if (!disablingCompaction.compareAndSet(false, true)) {
+                unsubscribeFuture.completeExceptionally(
+                        new SubscriptionBusyException("the subscription is 
deleting by another task"));
+                return;
             }
-            return ((CompactorSubscription) 
subscription).cleanCompactedLedger();
-        }).whenComplete((__, ex) -> {
-            if (ex != null) {
-                log.error("[{}][{}] Error cleaning compacted ledger", topic, 
subscriptionName, ex);
-                unsubscribeFuture.completeExceptionally(ex);
+        }
+        // Unsubscribe compaction cursor and delete compacted ledger.
+        currentCompaction.thenCompose(__ -> {
+            asyncDeleteCursor(subscriptionName, unsubscribeFuture);
+            return unsubscribeFuture;
+        }).thenAccept(__ -> {
+            try {
+                ((CompactorSubscription) subscription).cleanCompactedLedger();
+            } catch (Exception ex) {
+                Long compactedLedger = null;
+                Optional<CompactedTopicContext> compactedTopicContext = 
getCompactedTopicContext();
+                if (compactedTopicContext.isPresent() && 
compactedTopicContext.get().getLedger() != null) {
+                    compactedLedger = 
compactedTopicContext.get().getLedger().getId();
+                }
+                log.error("[{}][{}][{}] Error cleaning compacted ledger", 
topic, subscriptionName, compactedLedger, ex);
+            } finally {
+                // Reset the variable: disablingCompaction,
+                disablingCompaction.compareAndSet(true, false);
+            }
+        }).exceptionally(ex -> {
+            if (currentCompaction.isCompletedExceptionally()) {
+                log.warn("[{}][{}] Last compaction task failed", topic, 
subscriptionName);
             } else {
-                asyncDeleteCursor(subscriptionName, unsubscribeFuture);
+                log.warn("[{}][{}] Failed to delete cursor task failed", 
topic, subscriptionName);
             }
+            // Reset the variable: disablingCompaction,
+            disablingCompaction.compareAndSet(true, false);
+            unsubscribeFuture.completeExceptionally(ex);
+            return null;
         });
     }
 
@@ -3678,6 +3703,10 @@ public class PersistentTopic extends AbstractTopic 
implements Topic, AddEntryCal
                     log.info("[{}] Topic is closing or deleting, skip 
triggering compaction", topic);
                     return;
                 }
+                if (disablingCompaction.get()) {
+                    log.info("[{}] Compaction is disabling, skip triggering 
compaction", topic);
+                    return;
+                }
 
                 if (strategicCompactionMap.containsKey(topic)) {
                     currentCompaction = 
brokerService.pulsar().getStrategicCompactor()
diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java
index afd7b107ac8..035dc79db3f 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java
@@ -83,8 +83,15 @@ public class CompactedTopicImpl implements CompactedTopic {
             compactionHorizon = (PositionImpl) p;
 
             // delete the ledger from the old context once the new one is open
-            return compactedTopicContext.thenCompose(
-                    __ -> previousContext != null ? previousContext : 
CompletableFuture.completedFuture(null));
+            return compactedTopicContext.thenCompose(ctx -> {
+                if (ctx != null && ctx.getLedger() != null && 
ctx.getLedger().getId() == compactedLedgerId) {
+                    // Print an error log here, which is not expected.
+                    log.error("[__compaction] Using the same compacted ledger 
to override the old one, which is not"
+                                    + " expected and it may cause a ledger 
lost error. {} -> {}", compactedLedgerId,
+                            ctx.getLedger().getId());
+                }
+                return previousContext != null ? previousContext : 
CompletableFuture.completedFuture(null);
+            });
         }
     }
 
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/CompactionConcurrencyTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/CompactionConcurrencyTest.java
new file mode 100644
index 00000000000..7f7a32e6015
--- /dev/null
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/CompactionConcurrencyTest.java
@@ -0,0 +1,153 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pulsar.broker.service.persistent;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.fail;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+import org.apache.bookkeeper.mledger.Position;
+import org.apache.pulsar.broker.BrokerTestUtil;
+import org.apache.pulsar.client.admin.PulsarAdminException;
+import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.ProducerConsumerBase;
+import org.apache.pulsar.client.api.Schema;
+import org.apache.pulsar.common.naming.TopicName;
+import org.apache.pulsar.common.util.FutureUtil;
+import org.apache.pulsar.compaction.Compactor;
+import org.apache.zookeeper.MockZooKeeper;
+import org.awaitility.Awaitility;
+import org.awaitility.reflect.WhiteboxImpl;
+import org.testng.annotations.AfterClass;
+import org.testng.annotations.BeforeClass;
+import org.testng.annotations.Test;
+
+@Test(groups = "broker")
+public class CompactionConcurrencyTest extends ProducerConsumerBase {
+    // don't make this over 2000ms, otherwise the test will be flaky due to 
ZKSessionWatcher
+    static final int DELETE_OPERATION_DELAY_MS = 1900;
+
+    @BeforeClass
+    @Override
+    protected void setup() throws Exception {
+        super.internalSetup();
+        super.producerBaseSetup();
+    }
+
+    @AfterClass
+    @Override
+    protected void cleanup() throws Exception {
+        super.internalCleanup();
+    }
+
+    @Override
+    protected void doInitConf() throws Exception {
+        super.doInitConf();
+        // Disable the scheduled task: compaction.
+        
conf.setBrokerServiceCompactionMonitorIntervalInSeconds(Integer.MAX_VALUE);
+        // Disable the scheduled task: retention.
+        conf.setRetentionCheckIntervalInSeconds(Integer.MAX_VALUE);
+    }
+
+    private void triggerCompactionAndWait(String topicName) throws Exception {
+        PersistentTopic persistentTopic =
+                (PersistentTopic) 
pulsar.getBrokerService().getTopic(topicName, false).get().get();
+        persistentTopic.triggerCompaction();
+        Awaitility.await().untilAsserted(() -> {
+            Position lastConfirmPos = 
persistentTopic.getManagedLedger().getLastConfirmedEntry();
+            Position markDeletePos = persistentTopic
+                    
.getSubscription(Compactor.COMPACTION_SUBSCRIPTION).getCursor().getMarkDeletedPosition();
+            assertEquals(markDeletePos.getLedgerId(), 
lastConfirmPos.getLedgerId());
+            assertEquals(markDeletePos.getEntryId(), 
lastConfirmPos.getEntryId());
+        });
+    }
+
+    @Test
+    public void testDisableCompactionConcurrently() throws Exception {
+        String topicName = "persistent://public/default/" + 
BrokerTestUtil.newUniqueName("tp");
+        admin.topics().createNonPartitionedTopic(topicName);
+        admin.topicPolicies().setCompactionThreshold(topicName, 1);
+        admin.topics().createSubscription(topicName, "s1", MessageId.earliest);
+        var producer = 
pulsarClient.newProducer(Schema.STRING).topic(topicName).enableBatching(false).create();
+        producer.newMessage().key("k0").value("v0").send();
+        triggerCompactionAndWait(topicName);
+        admin.topics().deleteSubscription(topicName, "s1");
+        PersistentTopic persistentTopic =
+                (PersistentTopic) 
pulsar.getBrokerService().getTopic(topicName, false).get().get();
+        AtomicBoolean disablingCompaction = 
persistentTopic.disablingCompaction;
+
+        // Disable compaction.
+        // Inject a delay when the first time of deleting cursor.
+        AtomicInteger times = new AtomicInteger();
+        String cursorPath = String.format("/managed-ledgers/%s/__compaction",
+                TopicName.get(topicName).getPersistenceNamingEncoding());
+        admin.topicPolicies().removeCompactionThreshold(topicName);
+        mockZooKeeper.delay(DELETE_OPERATION_DELAY_MS, (op, path) -> {
+            return op == MockZooKeeper.Op.DELETE && cursorPath.equals(path) && 
times.incrementAndGet() == 1;
+        });
+        mockZooKeeperGlobal.delay(DELETE_OPERATION_DELAY_MS, (op, path) -> {
+            return op == MockZooKeeper.Op.DELETE && cursorPath.equals(path) && 
times.incrementAndGet() == 1;
+        });
+        AtomicReference<CompletableFuture<Void>> f1 = new 
AtomicReference<CompletableFuture<Void>>();
+        AtomicReference<CompletableFuture<Void>> f2 = new 
AtomicReference<CompletableFuture<Void>>();
+        new Thread(() -> {
+            f1.set(admin.topics().deleteSubscriptionAsync(topicName, 
"__compaction"));
+        }).start();
+        new Thread(() -> {
+            f2.set(admin.topics().deleteSubscriptionAsync(topicName, 
"__compaction"));
+        }).start();
+
+        // Verify: the next compaction will be skipped.
+        Awaitility.await().untilAsserted(() -> {
+            assertTrue(disablingCompaction.get());
+        });
+        producer.newMessage().key("k1").value("v1").send();
+        producer.newMessage().key("k2").value("v2").send();
+        CompletableFuture<Long> currentCompaction1 = 
persistentTopic.currentCompaction;
+                WhiteboxImpl.getInternalState(persistentTopic, 
"currentCompaction");
+        persistentTopic.triggerCompaction();
+        CompletableFuture<Long> currentCompaction2 = 
persistentTopic.currentCompaction;
+        assertTrue(currentCompaction1 == currentCompaction2);
+
+        // Verify: one of the requests should fail.
+        Awaitility.await().untilAsserted(() -> {
+            assertTrue(f1.get() != null);
+            assertTrue(f2.get() != null);
+            assertTrue(f1.get().isDone());
+            assertTrue(f2.get().isDone());
+            assertTrue(f1.get().isCompletedExceptionally() || 
f2.get().isCompletedExceptionally());
+            assertTrue(!f1.get().isCompletedExceptionally() || 
!f2.get().isCompletedExceptionally());
+        });
+        try {
+            f1.get().join();
+            f2.get().join();
+            fail("Should fail");
+        } catch (Exception ex) {
+            Throwable actEx = FutureUtil.unwrapCompletionException(ex);
+            assertTrue(actEx instanceof 
PulsarAdminException.PreconditionFailedException);
+        }
+
+        // cleanup.
+        producer.close();
+        admin.topics().delete(topicName, false);
+    }
+}
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/GetLastMessageIdCompactedTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/GetLastMessageIdCompactedTest.java
index 6c2d848bb7c..d3f0d97bf72 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/GetLastMessageIdCompactedTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/GetLastMessageIdCompactedTest.java
@@ -21,17 +21,20 @@ package org.apache.pulsar.compaction;
 import static org.testng.Assert.assertEquals;
 import static org.testng.Assert.assertFalse;
 import static org.testng.Assert.assertNotEquals;
-
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.fail;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
 import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl;
 import org.apache.bookkeeper.mledger.impl.PositionImpl;
 import org.apache.pulsar.broker.BrokerTestUtil;
 import org.apache.pulsar.broker.service.Topic;
 import org.apache.pulsar.broker.service.persistent.PersistentTopic;
+import org.apache.pulsar.client.admin.PulsarAdminException;
 import org.apache.pulsar.client.api.CompressionType;
 import org.apache.pulsar.client.api.Consumer;
 import org.apache.pulsar.client.api.Message;
@@ -39,11 +42,15 @@ import org.apache.pulsar.client.api.MessageId;
 import org.apache.pulsar.client.api.Producer;
 import org.apache.pulsar.client.api.ProducerBuilder;
 import org.apache.pulsar.client.api.ProducerConsumerBase;
+import org.apache.pulsar.client.api.Reader;
 import org.apache.pulsar.client.api.Schema;
 import org.apache.pulsar.client.impl.BatchMessageIdImpl;
 import org.apache.pulsar.client.impl.MessageIdImpl;
 import org.apache.pulsar.client.impl.ReaderImpl;
+import org.apache.pulsar.common.naming.TopicName;
 import org.apache.pulsar.common.util.FutureUtil;
+import org.apache.zookeeper.KeeperException;
+import org.apache.zookeeper.MockZooKeeper;
 import org.awaitility.Awaitility;
 import org.testng.annotations.AfterClass;
 import org.testng.annotations.BeforeClass;
@@ -309,6 +316,73 @@ public class GetLastMessageIdCompactedTest extends 
ProducerConsumerBase {
         admin.topics().delete(topicName, false);
     }
 
+    @DataProvider
+    public Object[][] isInjectedCursorDeleteError() {
+        return new Object[][] {
+                {false},
+                {true}
+        };
+    }
+
+    @Test(dataProvider = "isInjectedCursorDeleteError")
+    public void testReadMsgsAfterDisableCompaction(boolean 
isInjectedCursorDeleteError) throws Exception {
+        String topicName = "persistent://public/default/" + 
BrokerTestUtil.newUniqueName("tp");
+        admin.topics().createNonPartitionedTopic(topicName);
+        admin.topicPolicies().setCompactionThreshold(topicName, 1);
+        admin.topics().createSubscription(topicName, "s1", MessageId.earliest);
+        var producer = 
pulsarClient.newProducer(Schema.STRING).topic(topicName).enableBatching(false).create();
+        producer.newMessage().key("k0").value("v0").send();
+        producer.newMessage().key("k1").value("v1").send();
+        producer.newMessage().key("k2").value("v2").send();
+        triggerCompactionAndWait(topicName);
+        admin.topics().deleteSubscription(topicName, "s1");
+
+        // Disable compaction.
+        // Inject a failure that the first time to delete cursor will fail.
+        if (isInjectedCursorDeleteError) {
+            AtomicInteger times = new AtomicInteger();
+            String cursorPath = 
String.format("/managed-ledgers/%s/__compaction",
+                    TopicName.get(topicName).getPersistenceNamingEncoding());
+            admin.topicPolicies().removeCompactionThreshold(topicName);
+            mockZooKeeper.failConditional(KeeperException.Code.SESSIONEXPIRED, 
(op, path) -> {
+                return op == MockZooKeeper.Op.DELETE && 
cursorPath.equals(path) && times.incrementAndGet() == 1;
+            });
+            
mockZooKeeperGlobal.failConditional(KeeperException.Code.SESSIONEXPIRED, (op, 
path) -> {
+                return op == MockZooKeeper.Op.DELETE && 
cursorPath.equals(path) && times.incrementAndGet() == 1;
+            });
+            try {
+                admin.topics().deleteSubscription(topicName, "__compaction");
+                fail("Should fail");
+            } catch (Exception ex) {
+                assertTrue(ex instanceof 
PulsarAdminException.ServerSideErrorException);
+            }
+        }
+
+        // Create a reader with start at earliest.
+        // Verify: the reader will receive 3 messages.
+        admin.topics().unload(topicName);
+        Reader<String> reader = 
pulsarClient.newReader(Schema.STRING).topic(topicName).readCompacted(true)
+                .startMessageId(MessageId.earliest).create();
+        producer.newMessage().key("k3").value("v3").send();
+        assertTrue(reader.hasMessageAvailable());
+        Message<String> m0 = reader.readNext(10, TimeUnit.SECONDS);
+        assertEquals(m0.getValue(), "v0");
+        assertTrue(reader.hasMessageAvailable());
+        Message<String> m1 = reader.readNext(10, TimeUnit.SECONDS);
+        assertEquals(m1.getValue(), "v1");
+        assertTrue(reader.hasMessageAvailable());
+        Message<String> m2 = reader.readNext(10, TimeUnit.SECONDS);
+        assertEquals(m2.getValue(), "v2");
+        assertTrue(reader.hasMessageAvailable());
+        Message<String> m3 = reader.readNext(10, TimeUnit.SECONDS);
+        assertEquals(m3.getValue(), "v3");
+
+        // cleanup.
+        producer.close();
+        reader.close();
+        admin.topics().delete(topicName, false);
+    }
+
     @Test(dataProvider = "enabledBatch")
     public void testGetLastMessageIdAfterCompactionEndWithNullMsg(boolean 
enabledBatch) throws Exception {
         String topicName = "persistent://public/default/" + 
BrokerTestUtil.newUniqueName("tp");

Reply via email to