This is an automated email from the ASF dual-hosted git repository.
eolivelli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 478fd36 [Broker] waitingCursors potential heap memory leak (#13939)
478fd36 is described below
commit 478fd36227c2ede3e1162dd9a4361cffc5dbfceb
Author: gaozhangmin <[email protected]>
AuthorDate: Tue Feb 22 02:18:56 2022 +0800
[Broker] waitingCursors potential heap memory leak (#13939)
---
.../org/apache/bookkeeper/mledger/ManagedLedger.java | 7 +++++++
.../bookkeeper/mledger/impl/ManagedLedgerImpl.java | 4 ++++
.../service/persistent/PersistentSubscription.java | 1 +
.../pulsar/broker/admin/CreateSubscriptionTest.java | 18 ++++++++++++++++++
.../mledger/offload/jcloud/impl/MockManagedLedger.java | 5 +++++
5 files changed, 35 insertions(+)
diff --git
a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java
b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java
index 5c62dba..1f6e0d3 100644
---
a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java
+++
b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java
@@ -293,6 +293,13 @@ public interface ManagedLedger {
void deleteCursor(String name) throws InterruptedException,
ManagedLedgerException;
/**
+ * Remove a ManagedCursor from this ManagedLedger's waitingCursors.
+ *
+ * @param cursor the ManagedCursor
+ */
+ void removeWaitingCursor(ManagedCursor cursor);
+
+ /**
* Open a ManagedCursor asynchronously.
*
* @see #openCursor(String)
diff --git
a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
index 23fd63a..bfa3336 100644
---
a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
+++
b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
@@ -3485,6 +3485,10 @@ public class ManagedLedgerImpl implements ManagedLedger,
CreateCallback {
}
}
+ public void removeWaitingCursor(ManagedCursor cursor) {
+ this.waitingCursors.remove(cursor);
+ }
+
public boolean isCursorActive(ManagedCursor cursor) {
return activeCursors.get(cursor.getName()) != null;
}
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
index 6d74b53..fc4d0f2 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java
@@ -308,6 +308,7 @@ public class PersistentSubscription implements Subscription
{
if (dispatcher != null && dispatcher.getConsumers().isEmpty()) {
deactivateCursor();
+ topic.getManagedLedger().removeWaitingCursor(cursor);
if (!cursor.isDurable()) {
// If cursor is not durable, we need to clean up the
subscription as well
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/CreateSubscriptionTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/CreateSubscriptionTest.java
index 59ebdd5..5b43419 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/CreateSubscriptionTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/CreateSubscriptionTest.java
@@ -32,6 +32,7 @@ import java.util.concurrent.TimeUnit;
import java.io.IOException;
import javax.ws.rs.ClientErrorException;
import javax.ws.rs.core.Response.Status;
+import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl;
import org.apache.pulsar.broker.service.persistent.PersistentSubscription;
import org.apache.http.HttpResponse;
import org.apache.http.client.config.RequestConfig;
@@ -39,13 +40,16 @@ import org.apache.http.client.methods.HttpPut;
import org.apache.http.entity.StringEntity;
import org.apache.http.impl.client.CloseableHttpClient;
import org.apache.http.impl.client.HttpClientBuilder;
+import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.client.admin.PulsarAdminException;
import org.apache.pulsar.client.admin.PulsarAdminException.ConflictException;
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.Producer;
+import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.ProducerConsumerBase;
+import org.apache.pulsar.client.api.SubscriptionType;
import org.apache.pulsar.client.impl.MessageIdImpl;
import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.common.policies.data.PartitionedTopicStats;
@@ -348,4 +352,18 @@ public class CreateSubscriptionTest extends
ProducerConsumerBase {
producer.close();
}
+
+ @Test
+ public void testWaitingCurosrCausedMemoryLeak() throws Exception {
+ String topic = "persistent://my-property/my-ns/my-topic";
+ for (int i = 0; i < 10; i ++) {
+ Consumer<String> consumer =
pulsarClient.newConsumer(Schema.STRING).topic(topic)
+
.subscriptionType(SubscriptionType.Failover).subscriptionName("test" +
i).subscribe();
+ Awaitility.await().untilAsserted(() ->
assertTrue(consumer.isConnected()));
+ consumer.close();
+ }
+ PersistentTopic topicRef = (PersistentTopic)
pulsar.getBrokerService().getTopicReference(topic).get();
+ ManagedLedgerImpl ml =
(ManagedLedgerImpl)(topicRef.getManagedLedger());
+ assertEquals(ml.getWaitingCursorsCount(), 0);
+ }
}
diff --git
a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java
b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java
index d025cd1..c92c2f8 100644
---
a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java
+++
b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java
@@ -138,6 +138,11 @@ public class MockManagedLedger implements ManagedLedger {
}
@Override
+ public void removeWaitingCursor(ManagedCursor cursor) {
+
+ }
+
+ @Override
public void asyncOpenCursor(String name, AsyncCallbacks.OpenCursorCallback
callback, Object ctx) {
}