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");