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;
+ }
+ }
}