This is an automated email from the ASF dual-hosted git repository.
JackieTien97 pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new 82e70b270d6 [To dev/1.3] Enable show queries can be executed when the
available memory for Operators is insufficient
82e70b270d6 is described below
commit 82e70b270d64673a11206444b8440eea0111a9e4
Author: Weihao Li <[email protected]>
AuthorDate: Tue Jul 21 11:40:14 2026 +0800
[To dev/1.3] Enable show queries can be executed when the available memory
for Operators is insufficient
---
.../fragment/FragmentInstanceContext.java | 3 +
.../execution/operator/OperatorContext.java | 5 +
.../plan/planner/LocalExecutionPlanner.java | 108 ++++++++------
.../memory/FakedMemoryReservationManager.java | 3 +
.../planner/memory/MemoryReservationManager.java | 6 +
.../NotThreadSafeMemoryReservationManager.java | 113 +++++++++++---
.../memory/ThreadSafeMemoryReservationManager.java | 10 ++
.../LocalExecutionPlannerOperatorsMemoryTest.java | 164 +++++++++++++++++++++
8 files changed, 346 insertions(+), 66 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceContext.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceContext.java
index 35bd6237db8..70cc46428be 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceContext.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceContext.java
@@ -1224,6 +1224,9 @@ public class FragmentInstanceContext extends QueryContext
{
public void setHighestPriority(boolean highestPriority) {
this.highestPriority = highestPriority;
+ if (memoryReservationManager != null) {
+ memoryReservationManager.setHighestPriority(highestPriority);
+ }
}
public boolean isSingleSourcePath() {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorContext.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorContext.java
index 90bded8a0ee..cd49c736e02 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorContext.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/operator/OperatorContext.java
@@ -119,6 +119,11 @@ public class OperatorContext implements Accountable {
return getInstanceContext().getSessionInfo();
}
+ public boolean isHighestPriority() {
+ FragmentInstanceContext instanceContext = getInstanceContext();
+ return instanceContext != null && instanceContext.isHighestPriority();
+ }
+
public void recordScanAggregationFromRawDataCost(long costTimeInNanos) {
if (driverContext != null && driverContext.getFragmentInstanceContext() !=
null) {
driverContext
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java
index 97445414c65..b8aef336e10 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlanner.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.queryengine.plan.planner;
import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.queryengine.common.DeviceContext;
@@ -97,7 +98,7 @@ public class LocalExecutionPlanner {
context.invalidateParentPlanNodeIdToMemoryEstimator();
// check whether current free memory is enough to execute current query
- long estimatedMemorySize = checkMemory(memoryEstimator,
instanceContext.getStateMachine());
+ long estimatedMemorySize = checkMemory(memoryEstimator, instanceContext);
context.addPipelineDriverFactory(root, context.getDriverContext(),
estimatedMemorySize);
@@ -128,7 +129,7 @@ public class LocalExecutionPlanner {
context.invalidateParentPlanNodeIdToMemoryEstimator();
// check whether current free memory is enough to execute current query
- checkMemory(memoryEstimator, instanceContext.getStateMachine());
+ checkMemory(memoryEstimator, instanceContext);
context.addPipelineDriverFactory(root, context.getDriverContext(), 0);
@@ -139,7 +140,7 @@ public class LocalExecutionPlanner {
}
private long checkMemory(
- final PipelineMemoryEstimator memoryEstimator,
FragmentInstanceStateMachine stateMachine)
+ final PipelineMemoryEstimator memoryEstimator, FragmentInstanceContext
instanceContext)
throws MemoryNotEnoughException {
// if it is disabled, just return
@@ -152,43 +153,70 @@ public class LocalExecutionPlanner {
QueryRelatedResourceMetricSet.getInstance().updateEstimatedMemory(estimatedMemorySize);
+ long reservedBytes =
+ allocateOperatorsMemory(estimatedMemorySize,
instanceContext.isHighestPriority());
+ if (reservedBytes < 0) {
+ throw new MemoryNotEnoughException(
+ String.format(
+ "There is not enough memory to execute current fragment
instance, "
+ + "current remaining free memory is %dB, "
+ + "estimated memory usage for current fragment instance is
%dB",
+ freeMemoryForOperators, estimatedMemorySize));
+ }
+ FragmentInstanceStateMachine stateMachine =
instanceContext.getStateMachine();
+ if (reservedBytes > 0) {
+ stateMachine.addStateChangeListener(
+ newState -> {
+ if (newState.isDone()) {
+ try (SetThreadName fragmentInstanceName =
+ new
SetThreadName(stateMachine.getFragmentInstanceId().getFullId())) {
+ synchronized (this) {
+ this.freeMemoryForOperators += reservedBytes;
+ if (LOGGER.isDebugEnabled()) {
+ LOGGER.debug(
+ "[ReleaseMemory] release: {}, current remaining
memory: {}",
+ reservedBytes,
+ freeMemoryForOperators);
+ }
+ }
+ }
+ }
+ });
+ }
+ return reservedBytes;
+ }
+
+ /**
+ * Try to reserve bytes from the operators free-memory pool.
+ *
+ * @return allocated bytes on success ({@code > 0}), {@code 0} if nothing to
allocate or
+ * highest-priority fallback applies, {@code -1} if allocation failed
+ */
+ private long allocateOperatorsMemory(final long memoryInBytes, final boolean
isHighestPriority) {
+ if (memoryInBytes <= 0) {
+ return 0L;
+ }
synchronized (this) {
- if (estimatedMemorySize > freeMemoryForOperators) {
- throw new MemoryNotEnoughException(
- String.format(
- "There is not enough memory to execute current fragment
instance, "
- + "current remaining free memory is %dB, "
- + "estimated memory usage for current fragment instance is
%dB",
- freeMemoryForOperators, estimatedMemorySize));
- } else {
- freeMemoryForOperators -= estimatedMemorySize;
+ if (memoryInBytes <= freeMemoryForOperators) {
+ freeMemoryForOperators -= memoryInBytes;
if (LOGGER.isDebugEnabled()) {
LOGGER.debug(
"[ConsumeMemory] consume: {}, current remaining memory: {}",
- estimatedMemorySize,
+ memoryInBytes,
freeMemoryForOperators);
}
+ return memoryInBytes;
}
}
+ if (isHighestPriority) {
+ return 0L;
+ }
+ return -1L;
+ }
- stateMachine.addStateChangeListener(
- newState -> {
- if (newState.isDone()) {
- try (SetThreadName fragmentInstanceName =
- new
SetThreadName(stateMachine.getFragmentInstanceId().getFullId())) {
- synchronized (this) {
- this.freeMemoryForOperators += estimatedMemorySize;
- if (LOGGER.isDebugEnabled()) {
- LOGGER.debug(
- "[ReleaseMemory] release: {}, current remaining memory:
{}",
- estimatedMemorySize,
- freeMemoryForOperators);
- }
- }
- }
- }
- });
- return estimatedMemorySize;
+ @TestOnly
+ long allocateOperatorsMemoryForTest(final long memoryInBytes, final boolean
isHighestPriority) {
+ return allocateOperatorsMemory(memoryInBytes, isHighestPriority);
}
private QueryDataSourceType getQueryDataSourceType(DataDriverContext
dataDriverContext) {
@@ -239,16 +267,19 @@ public class LocalExecutionPlanner {
}
}
- public synchronized void reserveFromFreeMemoryForOperators(
+ public long reserveFromFreeMemoryForOperators(
final long memoryInBytes,
final long reservedBytes,
final String queryId,
- final String contextHolder) {
+ final String contextHolder,
+ final boolean isHighestPriority)
+ throws MemoryNotEnoughException {
if (memoryInBytes <= 0) {
throw new IllegalArgumentException(
"Bytes to reserve from free memory for operators should be larger
than 0");
}
- if (memoryInBytes > freeMemoryForOperators) {
+ long allocated = allocateOperatorsMemory(memoryInBytes, isHighestPriority);
+ if (allocated < 0) {
throw new MemoryNotEnoughException(
String.format(
"There is not enough memory for Query %s, the contextHolder is
%s,"
@@ -256,15 +287,8 @@ public class LocalExecutionPlanner {
+ "already reserved memory for this context in total is %dB,
"
+ "the memory requested this time is %dB",
queryId, contextHolder, freeMemoryForOperators, reservedBytes,
memoryInBytes));
- } else {
- freeMemoryForOperators -= memoryInBytes;
- if (LOGGER.isDebugEnabled()) {
- LOGGER.debug(
- "[ConsumeMemory] consume: {}, current remaining memory: {}",
- memoryInBytes,
- freeMemoryForOperators);
- }
}
+ return allocated;
}
public synchronized void releaseToFreeMemoryForOperators(final long
memoryInBytes) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/FakedMemoryReservationManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/FakedMemoryReservationManager.java
index 35ded4d6252..8d0c9ae5997 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/FakedMemoryReservationManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/FakedMemoryReservationManager.java
@@ -43,4 +43,7 @@ public class FakedMemoryReservationManager implements
MemoryReservationManager {
@Override
public void reserveMemoryVirtually(
final long bytesToBeReserved, final long bytesAlreadyReserved) {}
+
+ @Override
+ public void setHighestPriority(boolean isHighestPriority) {}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/MemoryReservationManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/MemoryReservationManager.java
index 62393673120..eddec15facc 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/MemoryReservationManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/MemoryReservationManager.java
@@ -72,4 +72,10 @@ public interface MemoryReservationManager {
* @param bytesAlreadyReserved the amount of memory that has already been
reserved
*/
void reserveMemoryVirtually(final long bytesToBeReserved, final long
bytesAlreadyReserved);
+
+ /**
+ * Mark this manager as highest-priority (e.g. SHOW QUERIES). When operators
memory is
+ * insufficient, allocation will fall back to zero bytes instead of failing.
+ */
+ void setHighestPriority(boolean isHighestPriority);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/NotThreadSafeMemoryReservationManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/NotThreadSafeMemoryReservationManager.java
index 4fa97f368ad..e4f211ea764 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/NotThreadSafeMemoryReservationManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/NotThreadSafeMemoryReservationManager.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.queryengine.plan.planner.memory;
+import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.db.queryengine.common.QueryId;
import org.apache.iotdb.db.queryengine.plan.planner.LocalExecutionPlanner;
@@ -26,6 +27,8 @@ import org.apache.tsfile.utils.Pair;
import javax.annotation.concurrent.NotThreadSafe;
+import static com.google.common.base.Preconditions.checkState;
+
@NotThreadSafe
public class NotThreadSafeMemoryReservationManager implements
MemoryReservationManager {
// To avoid reserving memory too frequently, we choose to do it in batches.
This is the lower
@@ -38,8 +41,16 @@ public class NotThreadSafeMemoryReservationManager
implements MemoryReservationM
private final String contextHolder;
+ private boolean isHighestPriority;
+
private long reservedBytesInTotal = 0;
+ /**
+ * Bytes logically reserved but not taken from the operators pool due to
highest-priority
+ * fallback.
+ */
+ private long fallbackBytesInTotal = 0;
+
private long bytesToBeReserved = 0;
private long bytesToBeReleased = 0;
@@ -49,6 +60,21 @@ public class NotThreadSafeMemoryReservationManager
implements MemoryReservationM
this.contextHolder = contextHolder;
}
+ @Override
+ public void setHighestPriority(boolean isHighestPriority) {
+ this.isHighestPriority = isHighestPriority;
+ }
+
+ @TestOnly
+ public long getReservedBytesInTotalForTest() {
+ return reservedBytesInTotal;
+ }
+
+ @TestOnly
+ public long getFallbackBytesInTotalForTest() {
+ return fallbackBytesInTotal;
+ }
+
@Override
public void reserveMemoryCumulatively(final long size) {
bytesToBeReserved += size;
@@ -60,52 +86,91 @@ public class NotThreadSafeMemoryReservationManager
implements MemoryReservationM
@Override
public void reserveMemoryImmediately() {
if (bytesToBeReserved != 0) {
- LOCAL_EXECUTION_PLANNER.reserveFromFreeMemoryForOperators(
- bytesToBeReserved, reservedBytesInTotal, queryId.getId(),
contextHolder);
- reservedBytesInTotal += bytesToBeReserved;
+ long actualReserved =
+ LOCAL_EXECUTION_PLANNER.reserveFromFreeMemoryForOperators(
+ bytesToBeReserved,
+ reservedBytesInTotal,
+ queryId.getId(),
+ contextHolder,
+ isHighestPriority);
+ if (actualReserved == 0) {
+ fallbackBytesInTotal += bytesToBeReserved;
+ } else {
+ reservedBytesInTotal += actualReserved;
+ }
bytesToBeReserved = 0;
}
}
+ public void reserveMemoryImmediately(final long size) {
+ if (size != 0) {
+ long actualReserved =
+ LOCAL_EXECUTION_PLANNER.reserveFromFreeMemoryForOperators(
+ size, reservedBytesInTotal, queryId.getId(), contextHolder,
isHighestPriority);
+ if (actualReserved == 0) {
+ fallbackBytesInTotal += size;
+ } else {
+ reservedBytesInTotal += actualReserved;
+ }
+ }
+ }
+
@Override
public void releaseMemoryCumulatively(final long size) {
+ if (size <= 0) {
+ return;
+ }
bytesToBeReleased += size;
if (bytesToBeReleased >= MEMORY_BATCH_THRESHOLD) {
- long bytesToRelease;
- if (bytesToBeReleased <= bytesToBeReserved) {
- bytesToBeReserved -= bytesToBeReleased;
- } else {
- bytesToRelease = bytesToBeReleased - bytesToBeReserved;
- bytesToBeReserved = 0;
-
LOCAL_EXECUTION_PLANNER.releaseToFreeMemoryForOperators(bytesToRelease);
- reservedBytesInTotal -= bytesToRelease;
- }
+ releaseBytesImmediately(bytesToBeReleased);
bytesToBeReleased = 0;
}
}
+ private void releaseBytesImmediately(final long size) {
+ long poolBytes = deductReleaseAccounting(size);
+ if (poolBytes > 0) {
+ LOCAL_EXECUTION_PLANNER.releaseToFreeMemoryForOperators(poolBytes);
+ }
+ }
+
+ /** Deduct release size from pending reserve, fallback quota, then pool
reservation in order. */
+ private long deductReleaseAccounting(final long size) {
+ long remaining = size;
+ if (remaining <= bytesToBeReserved) {
+ bytesToBeReserved -= remaining;
+ return 0L;
+ }
+ remaining -= bytesToBeReserved;
+ bytesToBeReserved = 0;
+
+ if (remaining <= fallbackBytesInTotal) {
+ fallbackBytesInTotal -= remaining;
+ return 0L;
+ }
+ remaining -= fallbackBytesInTotal;
+ fallbackBytesInTotal = 0;
+
+ reservedBytesInTotal -= remaining;
+ checkState(reservedBytesInTotal >= 0, "Released bytes has been larger than
reserved!");
+ return remaining;
+ }
+
@Override
public void releaseAllReservedMemory() {
if (reservedBytesInTotal != 0) {
LOCAL_EXECUTION_PLANNER.releaseToFreeMemoryForOperators(reservedBytesInTotal);
reservedBytesInTotal = 0;
- bytesToBeReserved = 0;
- bytesToBeReleased = 0;
}
+ fallbackBytesInTotal = 0;
+ bytesToBeReserved = 0;
+ bytesToBeReleased = 0;
}
@Override
public Pair<Long, Long> releaseMemoryVirtually(final long size) {
- if (bytesToBeReserved >= size) {
- bytesToBeReserved -= size;
- return new Pair<>(size, 0L);
- } else {
- long releasedBytesInReserved = bytesToBeReserved;
- long releasedBytesInTotal = size - bytesToBeReserved;
- bytesToBeReserved = 0;
- reservedBytesInTotal -= releasedBytesInTotal;
- return new Pair<>(releasedBytesInReserved, releasedBytesInTotal);
- }
+ long poolBytes = deductReleaseAccounting(size);
+ return new Pair<>(size - poolBytes, poolBytes);
}
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/ThreadSafeMemoryReservationManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/ThreadSafeMemoryReservationManager.java
index d167eae354f..0a1c6eee418 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/ThreadSafeMemoryReservationManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/ThreadSafeMemoryReservationManager.java
@@ -41,6 +41,11 @@ public class ThreadSafeMemoryReservationManager extends
NotThreadSafeMemoryReser
super.reserveMemoryImmediately();
}
+ @Override
+ public synchronized void reserveMemoryImmediately(final long size) {
+ super.reserveMemoryImmediately(size);
+ }
+
@Override
public synchronized void releaseMemoryCumulatively(long size) {
super.releaseMemoryCumulatively(size);
@@ -61,4 +66,9 @@ public class ThreadSafeMemoryReservationManager extends
NotThreadSafeMemoryReser
final long bytesToBeReserved, final long bytesAlreadyReserved) {
super.reserveMemoryVirtually(bytesToBeReserved, bytesAlreadyReserved);
}
+
+ @Override
+ public synchronized void setHighestPriority(boolean isHighestPriority) {
+ super.setHighestPriority(isHighestPriority);
+ }
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlannerOperatorsMemoryTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlannerOperatorsMemoryTest.java
new file mode 100644
index 00000000000..6d0cabb0443
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/LocalExecutionPlannerOperatorsMemoryTest.java
@@ -0,0 +1,164 @@
+/*
+ * 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.db.queryengine.plan.planner;
+
+import org.apache.iotdb.db.queryengine.common.QueryId;
+import
org.apache.iotdb.db.queryengine.plan.planner.memory.NotThreadSafeMemoryReservationManager;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Test;
+
+public class LocalExecutionPlannerOperatorsMemoryTest {
+
+ private static final LocalExecutionPlanner PLANNER =
LocalExecutionPlanner.getInstance();
+
+ private long bytesHeldByTest = 0L;
+
+ @After
+ public void tearDown() {
+ if (bytesHeldByTest > 0) {
+ PLANNER.releaseToFreeMemoryForOperators(bytesHeldByTest);
+ bytesHeldByTest = 0L;
+ }
+ }
+
+ @Test
+ public void testAllocateOperatorsMemoryFailsWhenInsufficient() {
+ long free = PLANNER.getFreeMemoryForOperators();
+ Assert.assertEquals(-1L, PLANNER.allocateOperatorsMemoryForTest(free +
1024L, false));
+ }
+
+ @Test
+ public void testAllocateOperatorsMemorySucceedsWhenAvailable() {
+ long request = Math.min(1024L, PLANNER.getFreeMemoryForOperators());
+ long reserved = PLANNER.allocateOperatorsMemoryForTest(request, false);
+ Assert.assertEquals(request, reserved);
+ bytesHeldByTest = reserved;
+ }
+
+ @Test
+ public void testHighestPriorityAllocatesWhenPoolHasRoom() {
+ long request = Math.min(1024L, PLANNER.getFreeMemoryForOperators());
+ if (request <= 0) {
+ return;
+ }
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+
+ long reserved = PLANNER.allocateOperatorsMemoryForTest(request, true);
+ Assert.assertEquals(request, reserved);
+ Assert.assertEquals(freeBefore - request,
PLANNER.getFreeMemoryForOperators());
+ bytesHeldByTest = reserved;
+ }
+
+ @Test
+ public void testHighestPriorityFallbackWhenPoolInsufficient() {
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+ long request = freeBefore + 1024L;
+
+ Assert.assertEquals(0L, PLANNER.allocateOperatorsMemoryForTest(request,
true));
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
+
+ @Test
+ public void
testMemoryReservationManagerHighestPriorityAllocatesWhenPoolHasRoom() {
+ long request = Math.min(1024L, PLANNER.getFreeMemoryForOperators());
+ if (request <= 0) {
+ return;
+ }
+
+ NotThreadSafeMemoryReservationManager manager =
+ new NotThreadSafeMemoryReservationManager(new QueryId("show_queries"),
"test");
+ manager.setHighestPriority(true);
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+
+ manager.reserveMemoryImmediately(request);
+ Assert.assertEquals(request, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore - request,
PLANNER.getFreeMemoryForOperators());
+
+ manager.releaseAllReservedMemory();
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
+
+ @Test
+ public void
testMemoryReservationManagerHighestPriorityFallbackWhenPoolInsufficient() {
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+ long request = freeBefore + 1024L;
+
+ NotThreadSafeMemoryReservationManager manager =
+ new NotThreadSafeMemoryReservationManager(new QueryId("show_queries"),
"test");
+ manager.setHighestPriority(true);
+
+ manager.reserveMemoryImmediately(request);
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(request, manager.getFallbackBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+
+ manager.releaseMemoryCumulatively(request);
+ Assert.assertEquals(0L, manager.getFallbackBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+
+ manager.releaseAllReservedMemory();
+ Assert.assertEquals(0L, manager.getFallbackBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
+
+ @Test
+ public void
testMemoryReservationManagerHighestPriorityFallbackReleaseViaBatchThreshold() {
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+ long request = freeBefore + MEMORY_BATCH_THRESHOLD;
+
+ NotThreadSafeMemoryReservationManager manager =
+ new NotThreadSafeMemoryReservationManager(new QueryId("show_queries"),
"test");
+ manager.setHighestPriority(true);
+
+ manager.reserveMemoryImmediately(request);
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(request, manager.getFallbackBytesInTotalForTest());
+
+ manager.releaseMemoryCumulatively(request);
+ Assert.assertEquals(0L, manager.getFallbackBytesInTotalForTest());
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
+
+ private static final long MEMORY_BATCH_THRESHOLD = 1024L * 1024L;
+
+ @Test
+ public void testMemoryReservationManagerNormalPriorityReserveAndRelease() {
+ long request = Math.min(1024L, PLANNER.getFreeMemoryForOperators());
+ if (request <= 0) {
+ return;
+ }
+
+ NotThreadSafeMemoryReservationManager manager =
+ new NotThreadSafeMemoryReservationManager(new QueryId("normal_query"),
"test");
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+
+ manager.reserveMemoryImmediately(request);
+ Assert.assertEquals(request, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore - request,
PLANNER.getFreeMemoryForOperators());
+
+ manager.releaseAllReservedMemory();
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
+}