This is an automated email from the ASF dual-hosted git repository. Wei-hao-Li pushed a commit to branch 1.3-2 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit da3f83454ddd499dacc515ed62b1379bbf9b3993 Author: Weihao Li <[email protected]> AuthorDate: Tue Jun 30 14:07:27 2026 +0800 Enable show queries can be executed when the available memory for Operators is insufficient (#18052) (cherry picked from commit a76341d940880b0f11102c09bf8a10183315cd94) --- .../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()); + } +}
