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

jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new ab7a911fdb2 [Subscription] Fix pending consensus request memory leak 
(#18429)
ab7a911fdb2 is described below

commit ab7a911fdb2d794dfee8465093f02217aafa0dd0
Author: Caideyipi <[email protected]>
AuthorDate: Thu Aug 13 15:08:55 2026 +0800

    [Subscription] Fix pending consensus request memory leak (#18429)
    
    * fix(subscription): bound pending consensus request memory
    
    * [Subscription] Release raw consensus requests without subscribers
    
    * [Subscription] Serialize before realtime queue admission
    
    * [Subscription] Keep queue absence diagnostic accurate
---
 .../common/request/IndexedConsensusRequest.java    |  44 ++++++++-
 .../consensus/iot/IoTConsensusServerImpl.java      |  11 ++-
 .../logdispatcher/IoTConsensusMemoryManager.java   |  70 +++++++------
 .../consensus/iot/logdispatcher/LogDispatcher.java |  11 +++
 .../subscription/SubscriptionQueueRegistry.java    |  10 +-
 .../IoTConsensusMemoryManagerTest.java             |  78 +++++++++++++++
 .../SubscriptionQueueRegistryTest.java             | 108 +++++++++++++++++++++
 .../consensus/ConsensusPrefetchingQueue.java       |  53 +++++++++-
 .../consensus/ConsensusPrefetchingQueueTest.java   |  92 +++++++++++++++++-
 9 files changed, 428 insertions(+), 49 deletions(-)

diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
index e90e39e1d12..834a752be6a 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java
@@ -47,11 +47,13 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
   private final List<IConsensusRequest> requests;
   private final List<ByteBuffer> serializedRequests;
   private long memorySize = 0;
-  private AtomicLong referenceCnt = new AtomicLong();
+  private long retainedMemorySize = 0;
+  private boolean serializedRequestsBuilt = false;
+  private final AtomicLong referenceCnt = new AtomicLong();
 
   public IndexedConsensusRequest(long searchIndex, List<IConsensusRequest> 
requests) {
     this.searchIndex = searchIndex;
-    this.requests = requests;
+    this.requests = new ArrayList<>(requests);
     this.syncIndex = -1L;
     this.serializedRequests = new ArrayList<>(requests.size());
   }
@@ -59,18 +61,23 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
   public IndexedConsensusRequest(
       long searchIndex, long syncIndex, List<IConsensusRequest> requests) {
     this.searchIndex = searchIndex;
-    this.requests = requests;
+    this.requests = new ArrayList<>(requests);
     this.syncIndex = syncIndex;
     this.serializedRequests = new ArrayList<>(requests.size());
   }
 
-  public void buildSerializedRequests() {
+  public synchronized void buildSerializedRequests() {
+    if (serializedRequestsBuilt) {
+      return;
+    }
     this.requests.forEach(
         r -> {
           ByteBuffer buffer = r.serializeToByteBuffer();
           this.serializedRequests.add(buffer);
           this.memorySize += Long.max(buffer.capacity(), r.getMemorySize());
+          this.retainedMemorySize += buffer.capacity() + r.getMemorySize();
         });
+    serializedRequestsBuilt = true;
   }
 
   @Override
@@ -90,6 +97,35 @@ public class IndexedConsensusRequest implements 
IConsensusRequest {
     return memorySize;
   }
 
+  /**
+   * Returns the memory retained while this request object is alive.
+   *
+   * <p>Before replication serialization, the request only retains its 
original request objects.
+   * Afterwards it retains both the originals and the serialized buffers. 
Batch memory continues to
+   * use {@link #getMemorySize()} because a batch no longer owns the original 
requests.
+   */
+  public synchronized long getRetainedMemorySize() {
+    if (serializedRequestsBuilt) {
+      return retainedMemorySize;
+    }
+    return requests.stream().mapToLong(IConsensusRequest::getMemorySize).sum();
+  }
+
+  /**
+   * Releases the original request objects after their serialized buffers have 
been materialized.
+   * Replication only needs the serialized buffers, while subscription 
delivery still needs the
+   * original objects and therefore must not call this method.
+   */
+  public synchronized void clearRequests() {
+    if (requests.isEmpty()) {
+      return;
+    }
+    if (serializedRequestsBuilt) {
+      retainedMemorySize -= 
requests.stream().mapToLong(IConsensusRequest::getMemorySize).sum();
+    }
+    requests.clear();
+  }
+
   public long getSearchIndex() {
     return searchIndex;
   }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
index 66d917bebc2..4cb109f0eb1 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
@@ -314,14 +314,15 @@ public class IoTConsensusServerImpl {
         // So we need to use the lock to ensure the `offer()` and 
`incrementAndGet()` are
         // in one transaction.
         synchronized (searchIndex) {
-          logDispatcher.offer(indexedConsensusRequest);
           // Deliver to subscription queues for real-time in-memory 
consumption.
           // Offer AFTER stateMachine.write() so that InsertNode has inferred 
types
           // and properly typed values (same timing as LogDispatcher).
-          final int sqCount = subscriptionQueueRegistry.size();
-          if (sqCount > 0) {
-            subscriptionQueueRegistry.offer(indexedConsensusRequest);
-          } else if (logger.isDebugEnabled()
+          final boolean offeredToSubscription =
+              subscriptionQueueRegistry.offer(indexedConsensusRequest);
+          logDispatcher.offer(indexedConsensusRequest, offeredToSubscription);
+          if (!offeredToSubscription
+              && subscriptionQueueRegistry.isEmpty()
+              && logger.isDebugEnabled()
               && indexedConsensusRequest.getSearchIndex() % 50 == 0) {
             // Log periodically when no subscription queues are registered
             logger.debug(
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
index 9bc5486f7e4..161494a5fe8 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
@@ -44,32 +44,35 @@ public class IoTConsensusMemoryManager {
   }
 
   public boolean reserve(IndexedConsensusRequest request) {
-    long prevRef = request.incRef();
-    if (prevRef == 0) {
-      boolean reserved = reserve(request.getMemorySize(), true);
-      if (reserved) {
-        if (logger.isDebugEnabled()) {
-          logger.debug(
-              IoTConsensusMessages.RESERVING_BYTES_FOR_REQUEST_SUCCEEDS,
-              request.getMemorySize(),
-              request.getSearchIndex(),
-              memoryBlock.getUsedMemoryInBytes());
-        }
-      } else {
-        request.decRef();
-        if (logger.isDebugEnabled()) {
-          logger.debug(
-              IoTConsensusMessages.RESERVING_BYTES_FOR_REQUEST_FAILS,
-              request.getMemorySize(),
-              request.getSearchIndex(),
-              memoryBlock.getUsedMemoryInBytes());
+    synchronized (request) {
+      long prevRef = request.incRef();
+      if (prevRef == 0) {
+        final long retainedMemorySize = request.getRetainedMemorySize();
+        boolean reserved = reserve(retainedMemorySize, true);
+        if (reserved) {
+          if (logger.isDebugEnabled()) {
+            logger.debug(
+                IoTConsensusMessages.RESERVING_BYTES_FOR_REQUEST_SUCCEEDS,
+                retainedMemorySize,
+                request.getSearchIndex(),
+                memoryBlock.getUsedMemoryInBytes());
+          }
+        } else {
+          request.decRef();
+          if (logger.isDebugEnabled()) {
+            logger.debug(
+                IoTConsensusMessages.RESERVING_BYTES_FOR_REQUEST_FAILS,
+                retainedMemorySize,
+                request.getSearchIndex(),
+                memoryBlock.getUsedMemoryInBytes());
+          }
         }
+        return reserved;
+      } else if (logger.isDebugEnabled()) {
+        logger.debug(IoTConsensusMessages.SKIP_MEMORY_RESERVATION, 
request.getSearchIndex());
       }
-      return reserved;
-    } else if (logger.isDebugEnabled()) {
-      logger.debug(IoTConsensusMessages.SKIP_MEMORY_RESERVATION, 
request.getSearchIndex());
+      return true;
     }
-    return true;
   }
 
   public boolean reserve(Batch batch) {
@@ -108,15 +111,18 @@ public class IoTConsensusMemoryManager {
   }
 
   public void free(IndexedConsensusRequest request) {
-    long prevRef = request.decRef();
-    if (prevRef == 1) {
-      free(request.getMemorySize(), true);
-      if (logger.isDebugEnabled()) {
-        logger.debug(
-            IoTConsensusMessages.FREED_BYTES_FOR_REQUEST,
-            request.getMemorySize(),
-            request.getSearchIndex(),
-            memoryBlock.getUsedMemoryInBytes());
+    synchronized (request) {
+      long prevRef = request.decRef();
+      if (prevRef == 1) {
+        final long retainedMemorySize = request.getRetainedMemorySize();
+        free(retainedMemorySize, true);
+        if (logger.isDebugEnabled()) {
+          logger.debug(
+              IoTConsensusMessages.FREED_BYTES_FOR_REQUEST,
+              retainedMemorySize,
+              request.getSearchIndex(),
+              memoryBlock.getUsedMemoryInBytes());
+        }
       }
     }
   }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
index 1622c2052fa..6250b361e38 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
@@ -184,9 +184,16 @@ public class LogDispatcher {
   }
 
   public void offer(IndexedConsensusRequest request) {
+    offer(request, true);
+  }
+
+  public void offer(IndexedConsensusRequest request, boolean 
keepRequestsForSubscription) {
     // we don't need to serialize and offer request when replicaNum is 1.
     if (!threads.isEmpty()) {
       request.buildSerializedRequests();
+      if (!keepRequestsForSubscription) {
+        request.clearRequests();
+      }
       synchronized (this) {
         threads.forEach(
             thread -> {
@@ -204,6 +211,10 @@ public class LogDispatcher {
               }
             });
       }
+    } else if (!keepRequestsForSubscription) {
+      // A single-replica region has no dispatcher thread, but still does not 
need the raw request
+      // after the state machine has applied it when subscriptions are 
disabled.
+      request.clearRequests();
     }
   }
 
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
index ea6872b8788..a0a9422e127 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
@@ -73,12 +73,15 @@ public class SubscriptionQueueRegistry {
     return new ArrayList<>(queues.values());
   }
 
-  public synchronized void offer(final IndexedConsensusRequest 
indexedConsensusRequest) {
+  public synchronized boolean offer(final IndexedConsensusRequest 
indexedConsensusRequest) {
     final int queueCount = queues.size();
     if (queueCount <= 0) {
-      return;
+      return false;
     }
 
+    // Subscription queues reserve request memory, including the serialized 
buffers.
+    indexedConsensusRequest.buildSerializedRequests();
+
     if (LOGGER.isDebugEnabled()) {
       LOGGER.debug(
           IoTConsensusMessages
@@ -91,8 +94,10 @@ public class SubscriptionQueueRegistry {
               : 
indexedConsensusRequest.getRequests().get(0).getClass().getSimpleName());
     }
 
+    boolean offeredToAnyQueue = false;
     for (final BlockingQueue<IndexedConsensusRequest> queue : queues.keySet()) 
{
       final boolean offered = queue.offer(indexedConsensusRequest);
+      offeredToAnyQueue |= offered;
       if (LOGGER.isDebugEnabled()) {
         LOGGER.debug(
             
IoTConsensusMessages.LOG_OFFER_RESULT_ARG_QUEUESIZE_ARG_QUEUEREMAINING_ARG_7ADC84C2,
@@ -125,5 +130,6 @@ public class SubscriptionQueueRegistry {
         }
       }
     }
+    return offeredToAnyQueue;
   }
 }
diff --git 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
index 3d89943772e..9bddb298716 100644
--- 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
+++ 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
@@ -21,6 +21,7 @@ package org.apache.iotdb.consensus.iot.logdispatcher;
 
 import org.apache.iotdb.commons.memory.AtomicLongMemoryBlock;
 import org.apache.iotdb.commons.memory.IMemoryBlock;
+import org.apache.iotdb.commons.request.IConsensusRequest;
 import org.apache.iotdb.consensus.common.request.ByteBufferConsensusRequest;
 import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
 
@@ -47,10 +48,12 @@ public class IoTConsensusMemoryManagerTest {
     previousMemoryBlock = 
IoTConsensusMemoryManager.getInstance().getMemoryBlock();
     IoTConsensusMemoryManager.getInstance()
         .setMemoryBlock(new AtomicLongMemoryBlock("Test", null, 
memoryBlockSize));
+    IoTConsensusMemoryManager.getInstance().reset();
   }
 
   @After
   public void tearDown() throws Exception {
+    IoTConsensusMemoryManager.getInstance().reset();
     
IoTConsensusMemoryManager.getInstance().setMemoryBlock(previousMemoryBlock);
   }
 
@@ -64,6 +67,60 @@ public class IoTConsensusMemoryManagerTest {
     testReserveAndRelease(3);
   }
 
+  @Test
+  public void testRawAndSerializedMemoryAreBothReservedOnce() {
+    final IndexedConsensusRequest request =
+        new IndexedConsensusRequest(
+            1, Collections.singletonList(new SizedConsensusRequest(10, 20)));
+
+    assertEquals(10L, request.getRetainedMemorySize());
+    request.buildSerializedRequests();
+    request.buildSerializedRequests();
+    assertEquals(20L, request.getMemorySize());
+    assertEquals(30L, request.getRetainedMemorySize());
+    assertEquals(1, request.getSerializedRequests().size());
+    request.clearRequests();
+    assertTrue(request.getRequests().isEmpty());
+    assertEquals(20L, request.getRetainedMemorySize());
+
+    assertTrue(IoTConsensusMemoryManager.getInstance().reserve(request));
+    assertTrue(IoTConsensusMemoryManager.getInstance().reserve(request));
+    assertEquals(
+        20L, 
IoTConsensusMemoryManager.getInstance().getMemoryBlock().getUsedMemoryInBytes());
+
+    IoTConsensusMemoryManager.getInstance().free(request);
+    assertEquals(
+        20L, 
IoTConsensusMemoryManager.getInstance().getMemoryBlock().getUsedMemoryInBytes());
+    IoTConsensusMemoryManager.getInstance().free(request);
+    assertEquals(
+        0L, 
IoTConsensusMemoryManager.getInstance().getMemoryBlock().getUsedMemoryInBytes());
+  }
+
+  @Test
+  public void testUnserializedRequestReservesRawMemory() {
+    final IndexedConsensusRequest request =
+        new IndexedConsensusRequest(
+            1, Collections.singletonList(new SizedConsensusRequest(10, 20)));
+
+    assertTrue(IoTConsensusMemoryManager.getInstance().reserve(request));
+    assertEquals(
+        10L, 
IoTConsensusMemoryManager.getInstance().getMemoryBlock().getUsedMemoryInBytes());
+    IoTConsensusMemoryManager.getInstance().free(request);
+    assertEquals(
+        0L, 
IoTConsensusMemoryManager.getInstance().getMemoryBlock().getUsedMemoryInBytes());
+  }
+
+  @Test
+  public void testClearUnserializedRequest() {
+    final IndexedConsensusRequest request =
+        new IndexedConsensusRequest(
+            1, Collections.singletonList(new SizedConsensusRequest(10, 20)));
+
+    request.clearRequests();
+    assertTrue(request.getRequests().isEmpty());
+    assertEquals(0L, request.getRetainedMemorySize());
+  }
+
   private void testReserveAndRelease(int numReservation) {
     int allocationSize = 1;
     long allocatedSize = 0;
@@ -100,4 +157,25 @@ public class IoTConsensusMemoryManagerTest {
     }
     assertEquals(0, 
IoTConsensusMemoryManager.getInstance().getMemorySizeInByte());
   }
+
+  private static final class SizedConsensusRequest implements 
IConsensusRequest {
+
+    private final long rawMemorySize;
+    private final int serializedMemorySize;
+
+    private SizedConsensusRequest(final long rawMemorySize, final int 
serializedMemorySize) {
+      this.rawMemorySize = rawMemorySize;
+      this.serializedMemorySize = serializedMemorySize;
+    }
+
+    @Override
+    public ByteBuffer serializeToByteBuffer() {
+      return ByteBuffer.allocate(serializedMemorySize);
+    }
+
+    @Override
+    public long getMemorySize() {
+      return rawMemorySize;
+    }
+  }
 }
diff --git 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistryTest.java
 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistryTest.java
new file mode 100644
index 00000000000..e9d0f7f3b0f
--- /dev/null
+++ 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistryTest.java
@@ -0,0 +1,108 @@
+/*
+ * 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.iotdb.consensus.iot.subscription;
+
+import org.apache.iotdb.commons.request.IConsensusRequest;
+import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
+import org.apache.iotdb.consensus.iot.SubscriptionWalRetentionPolicy;
+
+import org.junit.Test;
+
+import java.nio.ByteBuffer;
+import java.util.Collections;
+import java.util.concurrent.ArrayBlockingQueue;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertTrue;
+
+public class SubscriptionQueueRegistryTest {
+
+  @Test
+  public void testOfferWithoutQueuesDoesNotSerializeRequest() {
+    final SubscriptionQueueRegistry registry = new 
SubscriptionQueueRegistry("test");
+    final IndexedConsensusRequest request = newRequest();
+
+    assertFalse(registry.offer(request));
+    assertTrue(request.getSerializedRequests().isEmpty());
+  }
+
+  @Test
+  public void testOfferSerializesRequestBeforeQueueAdmission() {
+    final SubscriptionQueueRegistry registry = new 
SubscriptionQueueRegistry("test");
+    final InspectingQueue queue = new InspectingQueue();
+    registry.register(
+        queue,
+        new SubscriptionWalRetentionPolicy(
+            "test",
+            SubscriptionWalRetentionPolicy.UNBOUNDED,
+            SubscriptionWalRetentionPolicy.UNBOUNDED));
+    final IndexedConsensusRequest request = newRequest();
+
+    assertTrue(registry.offer(request));
+    assertEquals(1, request.getSerializedRequests().size());
+    assertEquals(1, queue.getSerializedRequestCountAtOffer());
+    assertSame(request, queue.poll());
+  }
+
+  private static IndexedConsensusRequest newRequest() {
+    return new IndexedConsensusRequest(
+        1, Collections.singletonList(new 
ByteBufferConsensusRequest(ByteBuffer.allocate(1))));
+  }
+
+  private static final class ByteBufferConsensusRequest implements 
IConsensusRequest {
+
+    private final ByteBuffer buffer;
+
+    private ByteBufferConsensusRequest(final ByteBuffer buffer) {
+      this.buffer = buffer;
+    }
+
+    @Override
+    public ByteBuffer serializeToByteBuffer() {
+      return buffer;
+    }
+
+    @Override
+    public long getMemorySize() {
+      return buffer.capacity();
+    }
+  }
+
+  private static final class InspectingQueue extends 
ArrayBlockingQueue<IndexedConsensusRequest> {
+
+    private int serializedRequestCountAtOffer = -1;
+
+    private InspectingQueue() {
+      super(1);
+    }
+
+    @Override
+    public boolean offer(final IndexedConsensusRequest request) {
+      serializedRequestCountAtOffer = request.getSerializedRequests().size();
+      return super.offer(request);
+    }
+
+    private int getSerializedRequestCountAtOffer() {
+      return serializedRequestCountAtOffer;
+    }
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
index f535328cdca..2d2060e15f7 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
@@ -26,6 +26,7 @@ import 
org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
 import org.apache.iotdb.consensus.iot.IoTConsensusServerImpl;
 import org.apache.iotdb.consensus.iot.SubscriptionWalRetentionPolicy;
 import org.apache.iotdb.consensus.iot.log.ConsensusReqReader;
+import org.apache.iotdb.consensus.iot.logdispatcher.IoTConsensusMemoryManager;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodeMiscMessages;
 import org.apache.iotdb.db.i18n.DataNodePipeMessages;
@@ -379,6 +380,9 @@ public class ConsensusPrefetchingQueue {
 
     private final Runnable wakeupHook;
     private final BooleanSupplier admissionSupplier;
+    private final IoTConsensusMemoryManager requestMemoryManager =
+        IoTConsensusMemoryManager.getInstance();
+    private final AtomicLong retainedRequestBytes = new AtomicLong(0L);
 
     private WakeableIndexedConsensusQueue(
         final int capacity, final Runnable wakeupHook, final BooleanSupplier 
admissionSupplier) {
@@ -394,7 +398,20 @@ public class ConsensusPrefetchingQueue {
         if (!admissionSupplier.getAsBoolean()) {
           return false;
         }
-        offered = super.offer(request);
+        if (!requestMemoryManager.reserve(request)) {
+          return false;
+        }
+        try {
+          offered = super.offer(request);
+        } catch (final Throwable t) {
+          requestMemoryManager.free(request);
+          throw t;
+        }
+        if (offered) {
+          retainedRequestBytes.addAndGet(request.getRetainedMemorySize());
+        } else {
+          requestMemoryManager.free(request);
+        }
       }
       if (offered) {
         wakeupHook.run();
@@ -404,7 +421,23 @@ public class ConsensusPrefetchingQueue {
 
     @Override
     public synchronized void clear() {
-      super.clear();
+      IndexedConsensusRequest request;
+      while (Objects.nonNull(request = super.poll())) {
+        release(request);
+      }
+    }
+
+    private void release(final IndexedConsensusRequest request) {
+      retainedRequestBytes.addAndGet(-request.getRetainedMemorySize());
+      requestMemoryManager.free(request);
+    }
+
+    private void release(final List<IndexedConsensusRequest> requests) {
+      requests.forEach(this::release);
+    }
+
+    private long getRetainedRequestBytes() {
+      return retainedRequestBytes.get();
     }
   }
 
@@ -1391,9 +1424,14 @@ public class ConsensusPrefetchingQueue {
             nextExpectedSearchIndex.get(),
             prefetchingQueue.size());
 
-        final MaterializationResult batchResult =
-            accumulateFromPending(
-                batch, lingerBatch, observedSeekGeneration, maxTablets, 
maxBatchBytes);
+        final MaterializationResult batchResult;
+        try {
+          batchResult =
+              accumulateFromPending(
+                  batch, lingerBatch, observedSeekGeneration, maxTablets, 
maxBatchBytes);
+        } finally {
+          pendingEntries.release(batch);
+        }
         if (batchResult != MaterializationResult.SUCCESS) {
           if (batchResult == MaterializationResult.WAL_GAP) {
             return PrefetchRoundResult.rescheduleAfter(WAL_GAP_RETRY_SLEEP_MS);
@@ -3812,6 +3850,10 @@ public class ConsensusPrefetchingQueue {
     return retainedTabletBytes.get();
   }
 
+  public long getRetainedRequestBytes() {
+    return pendingEntries.getRetainedRequestBytes();
+  }
+
   public long getSubscriptionMemoryLimitInBytes() {
     return subscriptionMemoryManager.getTotalMemorySizeInBytes();
   }
@@ -3880,6 +3922,7 @@ public class ConsensusPrefetchingQueue {
     result.put("prefetchingQueueSize", 
String.valueOf(prefetchingQueue.size()));
     result.put("inFlightEventsSize", String.valueOf(inFlightEvents.size()));
     result.put("pendingEntriesSize", String.valueOf(pendingEntries.size()));
+    result.put("retainedRequestBytes", 
String.valueOf(getRetainedRequestBytes()));
     result.put("retainedTabletBytes", 
String.valueOf(retainedTabletBytes.get()));
     result.put("memoryBlockedEntryBytes", 
String.valueOf(memoryBlockedEntryBytes));
     result.put("realtimeAdmissionBlocked", 
String.valueOf(realtimeAdmissionBlocked.get()));
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
index 845a797f149..b1fd02cc4c5 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
@@ -21,11 +21,15 @@ package org.apache.iotdb.db.subscription.broker.consensus;
 
 import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.consensus.DataRegionId;
+import org.apache.iotdb.commons.memory.AtomicLongMemoryBlock;
+import org.apache.iotdb.commons.memory.IMemoryBlock;
+import org.apache.iotdb.commons.request.IConsensusRequest;
 import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
 import org.apache.iotdb.consensus.iot.IoTConsensusServerImpl;
 import org.apache.iotdb.consensus.iot.SubscriptionWalRetentionPolicy;
 import org.apache.iotdb.consensus.iot.WriterSafeFrontierTracker;
 import org.apache.iotdb.consensus.iot.log.ConsensusReqReader;
+import org.apache.iotdb.consensus.iot.logdispatcher.IoTConsensusMemoryManager;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.queryengine.plan.statement.StatementTestUtils;
 import org.apache.iotdb.db.storageengine.dataregion.wal.io.ProgressWALReader;
@@ -56,6 +60,7 @@ import java.lang.reflect.Constructor;
 import java.lang.reflect.Field;
 import java.lang.reflect.Method;
 import java.lang.reflect.Modifier;
+import java.nio.ByteBuffer;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.List;
@@ -188,6 +193,54 @@ public class ConsensusPrefetchingQueueTest {
     assertTrue(queue.isEmpty());
   }
 
+  @Test
+  @SuppressWarnings("unchecked")
+  public void 
testPendingQueueUsesSharedConsensusMemoryBudgetAndClearReleasesIt() throws 
Exception {
+    final Class<?> queueClass =
+        Class.forName(ConsensusPrefetchingQueue.class.getName() + 
"$WakeableIndexedConsensusQueue");
+    final Constructor<?> constructor =
+        queueClass.getDeclaredConstructor(int.class, Runnable.class, 
BooleanSupplier.class);
+    constructor.setAccessible(true);
+    final Method retainedBytesMethod = 
queueClass.getDeclaredMethod("getRetainedRequestBytes");
+    retainedBytesMethod.setAccessible(true);
+
+    final IoTConsensusMemoryManager memoryManager = 
IoTConsensusMemoryManager.getInstance();
+    final IMemoryBlock previousMemoryBlock = memoryManager.getMemoryBlock();
+    final double previousQueueRatio = 
memoryManager.getMaxMemoryRatioForQueue();
+    final AtomicLongMemoryBlock testMemoryBlock =
+        new AtomicLongMemoryBlock("SubscriptionPendingTest", null, 100L);
+    final BlockingQueue<IndexedConsensusRequest> firstQueue =
+        (BlockingQueue<IndexedConsensusRequest>)
+            constructor.newInstance(8, (Runnable) () -> {}, (BooleanSupplier) 
() -> true);
+    final BlockingQueue<IndexedConsensusRequest> secondQueue =
+        (BlockingQueue<IndexedConsensusRequest>)
+            constructor.newInstance(8, (Runnable) () -> {}, (BooleanSupplier) 
() -> true);
+
+    memoryManager.init(testMemoryBlock, 0.6);
+    try {
+      final IndexedConsensusRequest sharedRequest = createSizedRequest(1L, 
10L, 30);
+      assertTrue(firstQueue.offer(sharedRequest));
+      assertTrue(secondQueue.offer(sharedRequest));
+      assertEquals(40L, testMemoryBlock.getUsedMemoryInBytes());
+      assertEquals(40L, retainedBytesMethod.invoke(firstQueue));
+      assertEquals(40L, retainedBytesMethod.invoke(secondQueue));
+
+      assertFalse(firstQueue.offer(createSizedRequest(2L, 10L, 30)));
+      assertEquals(40L, testMemoryBlock.getUsedMemoryInBytes());
+
+      firstQueue.clear();
+      assertEquals(40L, testMemoryBlock.getUsedMemoryInBytes());
+      assertEquals(0L, retainedBytesMethod.invoke(firstQueue));
+      secondQueue.clear();
+      assertEquals(0L, testMemoryBlock.getUsedMemoryInBytes());
+      assertEquals(0L, retainedBytesMethod.invoke(secondQueue));
+    } finally {
+      firstQueue.clear();
+      secondQueue.clear();
+      memoryManager.init(previousMemoryBlock, previousQueueRatio);
+    }
+  }
+
   @Test
   public void testLagIncludesLingeringBatchUntilCommitted() throws Exception {
     final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
@@ -230,8 +283,13 @@ public class ConsensusPrefetchingQueueTest {
               .setNodeId(7);
 
       assertNull(queue.poll("consumer"));
-      pendingEntries(queue).offer(request);
+      assertTrue(pendingEntries(queue).offer(request));
+      assertEquals(request.getRetainedMemorySize(), 
queue.getRetainedRequestBytes());
+      assertEquals(
+          String.valueOf(request.getRetainedMemorySize()),
+          queue.coreReportMessage().get("retainedRequestBytes"));
       queue.drivePrefetchOnce();
+      assertEquals(0L, queue.getRetainedRequestBytes());
 
       assertEquals(0, queue.getPrefetchedEventCount());
       assertEquals(1L, queue.getLag());
@@ -1800,6 +1858,17 @@ public class ConsensusPrefetchingQueueTest {
         .setNodeId(7);
   }
 
+  private static IndexedConsensusRequest createSizedRequest(
+      final long searchIndex, final long rawMemorySize, final int 
serializedMemorySize) {
+    final IndexedConsensusRequest request =
+        new IndexedConsensusRequest(
+            searchIndex,
+            Collections.singletonList(
+                new SizedConsensusRequest(rawMemorySize, 
serializedMemorySize)));
+    request.buildSerializedRequests();
+    return request;
+  }
+
   private static ConsensusSubscriptionCommitManager newCommitManager(final 
File systemDir)
       throws Exception {
     
IoTDBDescriptor.getInstance().getConfig().setSystemDir(systemDir.getAbsolutePath());
@@ -1850,4 +1919,25 @@ public class ConsensusPrefetchingQueueTest {
       return new Pair<>(DEFAULT_SAFELY_DELETED_SEARCH_INDEX, 0L);
     }
   }
+
+  private static final class SizedConsensusRequest implements 
IConsensusRequest {
+
+    private final long rawMemorySize;
+    private final int serializedMemorySize;
+
+    private SizedConsensusRequest(final long rawMemorySize, final int 
serializedMemorySize) {
+      this.rawMemorySize = rawMemorySize;
+      this.serializedMemorySize = serializedMemorySize;
+    }
+
+    @Override
+    public ByteBuffer serializeToByteBuffer() {
+      return ByteBuffer.allocate(serializedMemorySize);
+    }
+
+    @Override
+    public long getMemorySize() {
+      return rawMemorySize;
+    }
+  }
 }

Reply via email to