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 60b18800b6a [To dev/1.3] clone partial columns of aligned tvlist for
query (#18394)
60b18800b6a is described below
commit 60b18800b6a892317354d46fe5b013dbc9bc0629
Author: shuwenwei <[email protected]>
AuthorDate: Wed Aug 5 20:46:06 2026 +0800
[To dev/1.3] clone partial columns of aligned tvlist for query (#18394)
---
.../fragment/FragmentInstanceContext.java | 50 ++-
.../memory/FakedMemoryReservationManager.java | 3 +
.../planner/memory/MemoryReservationManager.java | 6 +
.../NotThreadSafeMemoryReservationManager.java | 16 +-
.../memory/ThreadSafeMemoryReservationManager.java | 5 +
.../schemaregion/utils/ResourceByPathUtils.java | 335 +++++++++++-----
.../memtable/AbstractWritableMemChunk.java | 24 +-
.../memtable/AlignedReadOnlyMemChunk.java | 4 +-
.../dataregion/memtable/ReadOnlyMemChunk.java | 4 +-
.../db/utils/datastructure/AlignedTVList.java | 437 ++++++++++++++++++---
.../iotdb/db/utils/datastructure/TVList.java | 10 +
.../fragment/FragmentInstanceExecutionTest.java | 21 +
.../LocalExecutionPlannerOperatorsMemoryTest.java | 46 +++
.../utils/ResourceByPathUtilsTest.java | 109 +++++
.../dataregion/memtable/TsFileProcessorTest.java | 28 +-
.../db/utils/datastructure/AlignedTVListTest.java | 181 ++++++++-
16 files changed, 1102 insertions(+), 177 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 70cc46428be..24b233efcdf 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
@@ -69,8 +69,10 @@ import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Optional;
import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
@@ -162,6 +164,9 @@ public class FragmentInstanceContext extends QueryContext {
private long closedUnseqFileNum = 0;
private boolean highestPriority = false;
+ // accessed value columns on each referenced AlignedTVList.
+ private final Map<TVList, Set<Integer>> alignedTVListColumnAccessMap = new
ConcurrentHashMap<>();
+
public static FragmentInstanceContext createFragmentInstanceContext(
FragmentInstanceId id,
FragmentInstanceStateMachine stateMachine,
@@ -218,6 +223,48 @@ public class FragmentInstanceContext extends QueryContext {
this.queryDataSourceType = queryDataSourceType;
}
+ /**
+ * Record columns of the AlignedTVList accessed by the query. This method is
called from
+ * prepareTvListMapForQuery with tvList.lockQueryList() held. Even though
the HashSet inside
+ * alignedTVListColumnAccessMap is not thread-safe, the calling pattern
guarantees thread safety
+ * without requiring additional synchronization.
+ *
+ * @param tvList the TVList being accessed
+ * @param columnIndexList list of column indices being accessed
+ */
+ public void putAccessedColumns(TVList tvList, List<Integer> columnIndexList)
{
+ Set<Integer> accessedColumns =
+ alignedTVListColumnAccessMap.computeIfAbsent(tvList, ignored -> new
HashSet<>());
+ columnIndexList.stream()
+ .filter(Objects::nonNull)
+ .forEach(
+ index -> {
+ if (index >= 0) {
+ accessedColumns.add(index);
+ }
+ });
+ }
+
+ /** Remove column-access metadata for an unpublished TVList when clone
preparation fails. */
+ public void removeAccessedColumns(TVList tvList) {
+ alignedTVListColumnAccessMap.remove(tvList);
+ }
+
+ /**
+ * Get columns of the AlignedTVList accessed by the query. This method is
called from
+ * prepareTvListMapForQuery with tvList.lockQueryList() held, ensuring that
no other thread can
+ * change accessed columns for the same TVList concurrently.
+ *
+ * @param tvList the TVList being accessed
+ * @return set of column indices being accessed
+ */
+ public Set<Integer> getAccessedAlignedColumns(TVList tvList) {
+ Set<Integer> accessedColumns = alignedTVListColumnAccessMap.get(tvList);
+ return accessedColumns == null
+ ? Collections.emptySet()
+ : Collections.unmodifiableSet(accessedColumns);
+ }
+
@TestOnly
public static FragmentInstanceContext createFragmentInstanceContext(
FragmentInstanceId id, FragmentInstanceStateMachine stateMachine) {
@@ -897,12 +944,12 @@ public class FragmentInstanceContext extends QueryContext
{
*/
private void releaseTVListOwnedByQuery() {
for (TVList tvList : tvListSet) {
- long tvListRamSize = tvList.calculateRamSize().getRamSize();
tvList.lockQueryList();
Set<QueryContext> queryContextSet = tvList.getQueryContextSet();
try {
queryContextSet.remove(this);
if (tvList.getOwnerQuery() == this) {
+ long tvListRamSize = tvList.calculateRamSize().getRamSize();
if (tvList.getReservedMemoryBytes() != tvListRamSize) {
LOGGER.warn(
"Release TVList owned by query: allocate size {}, release size
{}",
@@ -980,6 +1027,7 @@ public class FragmentInstanceContext extends QueryContext {
// release TVList/AlignedTVList owned by current query
releaseTVListOwnedByQuery();
+ alignedTVListColumnAccessMap.clear();
fileModCache = null;
nonExistentModFiles = null;
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 8d0c9ae5997..1742a8070b3 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
@@ -32,6 +32,9 @@ public class FakedMemoryReservationManager implements
MemoryReservationManager {
@Override
public void releaseMemoryCumulatively(long size) {}
+ @Override
+ public void releaseMemoryImmediately(long size) {}
+
@Override
public void releaseAllReservedMemory() {}
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 eddec15facc..9a9036b1d97 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
@@ -40,6 +40,12 @@ public interface MemoryReservationManager {
*/
void releaseMemoryCumulatively(final long size);
+ /**
+ * Release the given size immediately. This is used to roll back a
reservation when the operation
+ * protected by that reservation fails before ownership is published.
+ */
+ void releaseMemoryImmediately(final long size);
+
/**
* Release all reserved memory immediately. Make sure this method is called
when the lifecycle of
* this manager ends, Or the memory to be released in the batch may not be
released correctly.
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 e4f211ea764..514a2935f6c 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
@@ -79,7 +79,14 @@ public class NotThreadSafeMemoryReservationManager
implements MemoryReservationM
public void reserveMemoryCumulatively(final long size) {
bytesToBeReserved += size;
if (bytesToBeReserved >= MEMORY_BATCH_THRESHOLD) {
- reserveMemoryImmediately();
+ try {
+ reserveMemoryImmediately();
+ } catch (RuntimeException | Error failure) {
+ // reserveMemoryImmediately can fail only while asking the planner for
memory, before it
+ // updates this manager's counters. Keep the caller-visible
reservation operation atomic.
+ bytesToBeReserved -= size;
+ throw failure;
+ }
}
}
@@ -127,6 +134,13 @@ public class NotThreadSafeMemoryReservationManager
implements MemoryReservationM
}
}
+ @Override
+ public void releaseMemoryImmediately(final long size) {
+ if (size > 0) {
+ releaseBytesImmediately(size);
+ }
+ }
+
private void releaseBytesImmediately(final long size) {
long poolBytes = deductReleaseAccounting(size);
if (poolBytes > 0) {
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 0a1c6eee418..71676e5b77f 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
@@ -51,6 +51,11 @@ public class ThreadSafeMemoryReservationManager extends
NotThreadSafeMemoryReser
super.releaseMemoryCumulatively(size);
}
+ @Override
+ public synchronized void releaseMemoryImmediately(long size) {
+ super.releaseMemoryImmediately(size);
+ }
+
@Override
public synchronized void releaseAllReservedMemory() {
super.releaseAllReservedMemory();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtils.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtils.java
index fa2f603d6fa..6a3f7f23ab6 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtils.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtils.java
@@ -38,6 +38,7 @@ import
org.apache.iotdb.db.storageengine.dataregion.memtable.ReadOnlyMemChunk;
import org.apache.iotdb.db.storageengine.dataregion.modification.Modification;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.utils.ModificationUtils;
+import org.apache.iotdb.db.utils.datastructure.AlignedTVList;
import org.apache.iotdb.db.utils.datastructure.TVList;
import org.apache.tsfile.enums.TSDataType;
@@ -66,6 +67,7 @@ import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import static org.apache.iotdb.commons.path.AlignedPath.VECTOR_PLACEHOLDER;
@@ -121,7 +123,8 @@ public abstract class ResourceByPathUtils {
QueryContext context,
IWritableMemChunk memChunk,
boolean isWorkMemTable,
- Filter globalTimeFilter) {
+ Filter globalTimeFilter,
+ List<Integer> columnIndexList) {
// should copy globalTimeFilter because GroupByMonthFilter is stateful
Filter copyTimeFilter = null;
if (globalTimeFilter != null) {
@@ -141,118 +144,246 @@ public abstract class ResourceByPathUtils {
"Flushing/Working MemTable - add current query context to
immutable TVList's query list");
tvList.getQueryContextSet().add(context);
tvListQueryMap.put(tvList, tvList.rowCount());
+ // columnIndexList is to track column-level access for AlignedTVList.
+ // For TVList (primitive time series), it remains null and column
tracking is not needed.
+ if (columnIndexList != null && context instanceof
FragmentInstanceContext) {
+ ((FragmentInstanceContext) context).putAccessedColumns(tvList,
columnIndexList);
+ }
} finally {
tvList.unlockQueryList();
}
}
- // mutable tvlist
- TVList list = memChunk.getWorkingTVList();
- TVList cloneList = null;
- TVList.RamInfo listRamInfo = list.calculateRamSize();
- list.lockQueryList();
- try {
- if (copyTimeFilter != null
- && !copyTimeFilter.satisfyStartEndTime(list.getMinTime(),
list.getMaxTime())) {
- return tvListQueryMap;
- }
+ TVList.RamInfo listRamInfo = null;
+
+ // calculateRamSize (synchronized method on TVList) was previously called
before
+ // lockQueryList to avoid deadlock concerns. For partial clone of
AlignedTVList, however
+ // calculateRamSize must now be called inside the lockQueryList section
because it depends on
+ // accessing columns on the AlignedTVList.
+ // This is safe because the lock ordering — queryListLock must always be
acquired before the
+ // TVList intrinsic lock (via synchronized methods like calculateRamSize,
clone). So no AB-BA
+ // deadlock is possible.
+ while (true) {
+ // The working TVList may be replaced by a concurrent query via
clone-and-swap
+ // (memChunk.setWorkingTVList(clone)). A queryListLock held on a
detached candidate does
+ // not protect the current working TVList, so after acquiring the lock,
re-verify it is
+ // still the current working list under the memChunk lock. If it was
replaced while
+ // waiting for candidate's queryListLock, retry with the current one.
+ final TVList candidate = memChunk.getWorkingTVList();
+ candidate.lockQueryList();
+ try {
+ synchronized (memChunk) {
+ if (memChunk.getWorkingTVList() != candidate) {
+ continue;
+ }
+ }
- if (!isWorkMemTable) {
- /*
- * 1. Q1 queries this TVList while it is still in the working memtable
and records a smaller
- * visible row count.
- * 2. Later writes append out-of-order rows to the same TVList, then
FLUSH moves the
- * memtable to the flushing list.
- * 3. Q2 queries the flushing memtable. If Q2 directly reuses the
original mutable TVList,
- * Q2's query-side sort may reorder the indices in place.
- * 4. Q1 continues to read with its old row count and the reordered
indices. The converted
- * value index can exceed Q1's bitmap range and cause out-of-bound
access.
- *
- * Therefore, this flushing branch can reuse the original list only
when it is already
- * sorted or no active query is using it. Otherwise, Q2 should read
from
- * workingListForFlush.
- */
- boolean canUseListDirectly = list.isSorted() ||
list.getQueryContextSet().isEmpty();
- LOGGER.debug(
- "Flushing MemTable - add current query context to mutable TVList's
query list");
- if (canUseListDirectly) {
- list.getQueryContextSet().add(context);
- tvListQueryMap.put(list, list.rowCount());
- } else {
- TVList workingListForFlushSort =
memChunk.initWorkingListForFlushIfNecessary(list, true);
+ if (copyTimeFilter != null
+ && !copyTimeFilter.satisfyStartEndTime(
+ candidate.getMinTime(), candidate.getMaxTime())) {
+ return tvListQueryMap;
+ }
+
+ if (!isWorkMemTable) {
/*
- * The query will read from workingListForFlushSort, but
cloneForFlushSort() only clones
- * times and indices. The value arrays and bitmaps are still shared
with the original
- * list.
- *
- * Therefore, this query must also hold the original list until it
finishes. Adding
- * context to list.getQueryContextSet() lets flush/query cleanup see
that the original
- * list is still in use. Adding list to context.tvListSet makes
- * releaseTVListOwnedByQuery() remove this context from the original
list later.
+ * 1. Q1 queries this TVList while it is still in the working
memtable and records a smaller
+ * visible row count.
+ * 2. Later writes append out-of-order rows to the same TVList, then
FLUSH moves the
+ * memtable to the flushing list.
+ * 3. Q2 queries the flushing memtable. If Q2 directly reuses the
original mutable TVList,
+ * Q2's query-side sort may reorder the indices in place.
+ * 4. Q1 continues to read with its old row count and the reordered
indices. The converted
+ * value index can exceed Q1's bitmap range and cause
out-of-bound access.
*
- * Do not put the original list into tvListQueryMap here. The actual
read path must use
- * workingListForFlushSort to avoid sorting the original list in
place.
+ * Therefore, this flushing branch can reuse the original list only
when it is already
+ * sorted or no active query is using it. Otherwise, Q2 should read
from
+ * workingListForFlush.
*/
- list.getQueryContextSet().add(context);
- context.addTVListToSet(Collections.singleton(list));
- workingListForFlushSort.getQueryContextSet().add(context);
- tvListQueryMap.put(workingListForFlushSort,
workingListForFlushSort.rowCount());
+ boolean canUseListDirectly =
+ candidate.isSorted() || candidate.getQueryContextSet().isEmpty();
+ LOGGER.debug(
+ "Flushing MemTable - add current query context to mutable
TVList's query list");
+ if (canUseListDirectly) {
+ candidate.getQueryContextSet().add(context);
+ tvListQueryMap.put(candidate, candidate.rowCount());
+ } else {
+ TVList workingListForFlushSort =
+ memChunk.initWorkingListForFlushIfNecessary(candidate, true);
+ /*
+ * The query will read from workingListForFlushSort, but
cloneForFlushSort() only clones
+ * times and indices. The value arrays and bitmaps are still
shared with the original
+ * list.
+ *
+ * Therefore, this query must also hold the original list until it
finishes. Adding
+ * context to list.getQueryContextSet() lets flush/query cleanup
see that the original
+ * list is still in use. Adding list to context.tvListSet makes
+ * releaseTVListOwnedByQuery() remove this context from the
original list later.
+ *
+ * Do not put the original list into tvListQueryMap here. The
actual read path must use
+ * workingListForFlushSort to avoid sorting the original list in
place.
+ */
+ candidate.getQueryContextSet().add(context);
+ context.addTVListToSet(Collections.singleton(candidate));
+ // Query preparation is serialized by candidate's query-list lock,
but cleanup removes
+ // the context under workingListForFlushSort's own lock. Use the
same lock for this add
+ // to avoid concurrently mutating its HashSet. The lock order here
is candidate first,
+ // then workingListForFlushSort; cleanup never holds both locks at
the same time.
+ workingListForFlushSort.lockQueryList();
+ try {
+ workingListForFlushSort.getQueryContextSet().add(context);
+ } finally {
+ workingListForFlushSort.unlockQueryList();
+ }
+ tvListQueryMap.put(workingListForFlushSort,
workingListForFlushSort.rowCount());
+ }
+
+ // columnIndexList is to track column-level access for AlignedTVList.
+ // For TVList (primitive time series), it remains null and column
tracking is not needed.
+ if (columnIndexList != null && context instanceof
FragmentInstanceContext) {
+ ((FragmentInstanceContext) context).putAccessedColumns(candidate,
columnIndexList);
+ }
+ return tvListQueryMap;
}
- } else {
- if (list.isSorted() || list.getQueryContextSet().isEmpty()) {
+
+ if (candidate.isSorted() || candidate.getQueryContextSet().isEmpty()) {
LOGGER.debug(
"Working MemTable - add current query context to mutable
TVList's query list when it's sorted or no other query on it");
- list.getQueryContextSet().add(context);
- tvListQueryMap.put(list, list.rowCount());
- } else {
- /*
- * +----------------------+
- * | MemTable |
- * | |
- * | +------------+ | +-----------------+
- * | | TVList |<---+--+ +---+ Previous Query |
- * | +-----^------+ | | | +-----------------+
- * | | | | |
- * +----------+-----------+ | | +----------------+
- * | Clone +---+---+ Current Query |
- * +-----+------+ | +----------------+
- * | TVList | <---------+
- * +------------+
- */
- LOGGER.debug(
- "Working MemTable - clone mutable TVList and replace old TVList
in working MemTable");
- QueryContext firstQuery =
list.getQueryContextSet().iterator().next();
- // reserve query memory
- if (firstQuery instanceof FragmentInstanceContext) {
- MemoryReservationManager memoryReservationManager =
- ((FragmentInstanceContext)
firstQuery).getMemoryReservationContext();
-
memoryReservationManager.reserveMemoryCumulatively(listRamInfo.getRamSize());
- list.setReservedMemoryBytes(listRamInfo.getRamSize());
+ candidate.getQueryContextSet().add(context);
+ tvListQueryMap.put(candidate, candidate.rowCount());
+
+ // columnIndexList is to track column-level access for AlignedTVList.
+ // For TVList (primitive time series), it remains null and column
tracking is not needed.
+ if (columnIndexList != null && context instanceof
FragmentInstanceContext) {
+ ((FragmentInstanceContext) context).putAccessedColumns(candidate,
columnIndexList);
+ }
+ return tvListQueryMap;
+ }
+
+ /*
+ * +----------------------+
+ * | MemTable |
+ * | |
+ * | +------------+ | +-----------------+
+ * | | TVList |<---+--+ +---+ Previous Query |
+ * | +-----^------+ | | | +-----------------+
+ * | | | | |
+ * +----------+-----------+ | | +----------------+
+ * | Clone +---+---+ Current Query |
+ * +-----+------+ | +----------------+
+ * | TVList | <---------+
+ * +------------+
+ */
+ LOGGER.debug(
+ "Working MemTable - clone mutable TVList and replace old TVList in
working MemTable");
+
+ synchronized (memChunk) {
+ // Re-check defensively before cloning and publishing the
replacement. The clone and the
+ // working-list swap must be done in the same memChunk critical
section, so a concurrent
+ // query can never observe a working TVList whose columns have
already been moved away.
+ if (memChunk.getWorkingTVList() != candidate) {
+ continue;
}
- list.setOwnerQuery(firstQuery);
- // clone TVList
- cloneList = list.clone();
- cloneList.getQueryContextSet().add(context);
- tvListQueryMap.put(cloneList, cloneList.rowCount());
+ // calculateRamSize (synchronized method on TVList) was previously
called before
+ // lockQueryList to avoid deadlock concerns. For partial clone of
AlignedTVList, however
+ // calculateRamSize must now be called inside the lockQueryList
section because it depends
+ // on accessing columns on the AlignedTVList.
+ // This is safe because the lock ordering - queryListLock must
always be acquired before
+ // the TVList intrinsic lock (via synchronized methods like
calculateRamSize, clone). So
+ // no AB-BA deadlock is possible.
+ Set<Integer> columnsToClone = candidate.getAccessedColumnsForQuery();
+ listRamInfo =
+ (columnsToClone == null)
+ ? candidate.calculateRamSize()
+ : ((AlignedTVList)
candidate).calculateRamSize(columnsToClone);
+
+ QueryContext firstQuery =
candidate.getQueryContextSet().iterator().next();
+ TVList cloneList = null;
+ AlignedTVList.PartialClonePlan partialClonePlan = null;
+ FragmentInstanceContext cloneContext =
+ columnIndexList != null && context instanceof
FragmentInstanceContext
+ ? (FragmentInstanceContext) context
+ : null;
+ MemoryReservationManager memoryReservationManager =
+ firstQuery instanceof FragmentInstanceContext
+ ? ((FragmentInstanceContext)
firstQuery).getMemoryReservationContext()
+ : null;
+ boolean reservationNeedsRollback = false;
+ boolean replacementPublished = false;
+ try {
+ // Reserve before allocating the clone, so this transient memory
increase is still
+ // protected by query-memory admission control. Ownership is not
published yet, and a
+ // later preparation failure rolls this exact reservation back
immediately.
+ if (memoryReservationManager != null) {
+
memoryReservationManager.reserveMemoryCumulatively(listRamInfo.getRamSize());
+ reservationNeedsRollback = true;
+ }
+
+ // Clone and validate without changing the source list.
PartialClonePlan.commit is the
+ // only destructive step and is allocation-free.
+ if (columnsToClone == null) {
+ cloneList = candidate.clone();
+ } else {
+ partialClonePlan = ((AlignedTVList)
candidate).preparePartialClone(columnsToClone);
+ cloneList = partialClonePlan.getCloneList();
+ }
+
+ cloneList.getQueryContextSet().add(context);
+ tvListQueryMap.put(cloneList, cloneList.rowCount());
+ if (cloneContext != null) {
+ cloneContext.putAccessedColumns(cloneList, columnIndexList);
+ }
+
+ if (partialClonePlan != null) {
+ partialClonePlan.commit();
+ }
+ memChunk.setWorkingTVList(cloneList);
+ replacementPublished = true;
+
+ // Publish query ownership only after the replacement is fully
committed. The
+ // candidate query-list lock prevents its owner from being
released concurrently.
+ if (memoryReservationManager != null) {
+ candidate.setReservedMemoryBytes(listRamInfo.getRamSize());
+ }
+ candidate.setOwnerQuery(firstQuery);
+ reservationNeedsRollback = false;
+ return tvListQueryMap;
+ } catch (RuntimeException | Error failure) {
+ if (reservationNeedsRollback) {
+ try {
+
memoryReservationManager.releaseMemoryImmediately(listRamInfo.getRamSize());
+ } catch (RuntimeException | Error rollbackFailure) {
+ failure.addSuppressed(rollbackFailure);
+ }
+ }
+
+ // Before commit, remove the only external reference installed for
the unpublished
+ // clone. Its arrays can then be reclaimed while candidate remains
the working list.
+ if (!replacementPublished && cloneList != null) {
+ cloneList.getQueryContextSet().remove(context);
+ tvListQueryMap.remove(cloneList);
+ if (cloneContext != null) {
+ cloneContext.removeAccessedColumns(cloneList);
+ }
+ }
+ throw failure;
+ }
+ }
+ } catch (MemoryNotEnoughException ex) {
+ if (listRamInfo != null) {
+ LOGGER.warn(
+ "Failed to reserve memory for TVList: ramSize {}, timestampsSize
{}, arrayMemCost {}, rowCount {}, dataTypes {}",
+ listRamInfo.getRamSize(),
+ listRamInfo.getTimestampsSize(),
+ listRamInfo.getArrayMemCost(),
+ listRamInfo.getRowCount(),
+ listRamInfo.getDataTypes());
}
+ throw ex;
+ } finally {
+ candidate.unlockQueryList();
}
- } catch (MemoryNotEnoughException ex) {
- LOGGER.warn(
- "Failed to reserve memory for TVList: ramSize {}, timestampsSize {},
arrayMemCost {}, rowCount {}, dataTypes {}",
- listRamInfo.getRamSize(),
- listRamInfo.getTimestampsSize(),
- listRamInfo.getArrayMemCost(),
- listRamInfo.getRowCount(),
- listRamInfo.getDataTypes());
- throw ex;
- } finally {
- list.unlockQueryList();
- }
- if (cloneList != null) {
- memChunk.setWorkingTVList(cloneList);
}
- return tvListQueryMap;
}
}
@@ -400,15 +531,15 @@ class AlignedResourceByPathUtils extends
ResourceByPathUtils {
return null;
}
- // prepare AlignedTVList for query. It should clone TVList if necessary.
- Map<TVList, Integer> alignedTvListQueryMap =
- prepareTvListMapForQuery(
- context, alignedMemChunk, modsToMemtable == null,
globalTimeFilter);
-
// column index list for the query
List<Integer> columnIndexList =
alignedMemChunk.buildColumnIndexList(partialPath.getSchemaList());
+ // prepare AlignedTVList for query. It should clone TVList if necessary.
+ Map<TVList, Integer> alignedTvListQueryMap =
+ prepareTvListMapForQuery(
+ context, alignedMemChunk, modsToMemtable == null,
globalTimeFilter, columnIndexList);
+
List<List<TimeRange>> deletionList = null;
if (modsToMemtable != null) {
deletionList =
@@ -571,7 +702,7 @@ class MeasurementResourceByPathUtils extends
ResourceByPathUtils {
memTableMap.get(deviceID).getMemChunkMap().get(partialPath.getMeasurement());
// prepare TVList for query. It should clone TVList if necessary.
Map<TVList, Integer> tvListQueryMap =
- prepareTvListMapForQuery(context, memChunk, modsToMemtable == null,
globalTimeFilter);
+ prepareTvListMapForQuery(context, memChunk, modsToMemtable == null,
globalTimeFilter, null);
List<TimeRange> deletionList = null;
if (modsToMemtable != null) {
deletionList =
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AbstractWritableMemChunk.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AbstractWritableMemChunk.java
index 6c773942fb7..0f5edd22adb 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AbstractWritableMemChunk.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AbstractWritableMemChunk.java
@@ -24,6 +24,7 @@ import
org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContex
import org.apache.iotdb.db.queryengine.execution.fragment.QueryContext;
import
org.apache.iotdb.db.queryengine.plan.planner.memory.MemoryReservationManager;
import
org.apache.iotdb.db.storageengine.dataregion.wal.buffer.IWALByteBufferView;
+import org.apache.iotdb.db.utils.datastructure.AlignedTVList;
import org.apache.iotdb.db.utils.datastructure.BatchEncodeInfo;
import org.apache.iotdb.db.utils.datastructure.TVList;
@@ -37,6 +38,7 @@ import org.slf4j.LoggerFactory;
import java.util.Iterator;
import java.util.List;
+import java.util.Set;
import java.util.concurrent.BlockingQueue;
public abstract class AbstractWritableMemChunk implements IWritableMemChunk {
@@ -100,19 +102,39 @@ public abstract class AbstractWritableMemChunk implements
IWritableMemChunk {
}
}
+ /**
+ * Try to release the TVList. If there are active queries, transfer memory
ownership to the first
+ * query. For AlignedTVList, this will release non-query columns before
transferring to reduce
+ * memory footprint.
+ */
private void tryReleaseTvList(TVList tvList) {
- long tvListRamSize = tvList.calculateRamSize().getRamSize();
tvList.lockQueryList();
try {
if (tvList.getQueryContextSet().isEmpty()) {
tvList.clear();
} else {
QueryContext firstQuery =
tvList.getQueryContextSet().iterator().next();
+
+ // For AlignedTVList with active queries, release non-query columns
before
+ // transferring memory ownership to reduce memory footprint.
+ if (tvList instanceof AlignedTVList) {
+ AlignedTVList alignedTVList = (AlignedTVList) tvList;
+
+ // Get the union of all columns accessed by queries
+ Set<Integer> accessedColumns =
alignedTVList.getAccessedColumnsForQuery();
+
+ if (accessedColumns != null && !accessedColumns.isEmpty()) {
+ // Release non-query columns to reduce memory before ownership
transfer
+ alignedTVList.releaseNonQueryColumns(accessedColumns);
+ }
+ }
+
// transfer memory from write process to read process. Here it
reserves read memory and
// releaseFlushedMemTable will release write memory.
if (firstQuery instanceof FragmentInstanceContext) {
MemoryReservationManager memoryReservationManager =
((FragmentInstanceContext)
firstQuery).getMemoryReservationContext();
+ long tvListRamSize = tvList.calculateRamSize().getRamSize();
memoryReservationManager.reserveMemoryCumulatively(tvListRamSize);
tvList.setReservedMemoryBytes(tvListRamSize);
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java
index bb2ee311d30..f54af9cfcbb 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/AlignedReadOnlyMemChunk.java
@@ -122,12 +122,12 @@ public class AlignedReadOnlyMemChunk extends
ReadOnlyMemChunk {
// We must update queryRowCount here, otherwise, it may be used later
to build
// BitMaps, causing bitmap array size mismatch and possible out of
bound.
entry.setValue(alignedTvList.sort());
- long alignedTvListRamSize =
alignedTvList.calculateRamSize().getRamSize();
alignedTvList.lockQueryList();
try {
FragmentInstanceContext ownerQuery =
(FragmentInstanceContext) alignedTvList.getOwnerQuery();
if (ownerQuery != null) {
+ long alignedTvListRamSize =
alignedTvList.calculateRamSize().getRamSize();
long deltaBytes = alignedTvListRamSize -
alignedTvList.getReservedMemoryBytes();
if (deltaBytes > 0) {
ownerQuery.getMemoryReservationContext().reserveMemoryCumulatively(deltaBytes);
@@ -367,12 +367,12 @@ public class AlignedReadOnlyMemChunk extends
ReadOnlyMemChunk {
int queryLength = entry.getValue();
if (!alignedTvList.isSorted() && queryLength >
alignedTvList.seqRowCount()) {
entry.setValue(alignedTvList.sort());
- long alignedTvListRamSize =
alignedTvList.calculateRamSize().getRamSize();
alignedTvList.lockQueryList();
try {
FragmentInstanceContext ownerQuery =
(FragmentInstanceContext) alignedTvList.getOwnerQuery();
if (ownerQuery != null) {
+ long alignedTvListRamSize =
alignedTvList.calculateRamSize().getRamSize();
long deltaBytes = alignedTvListRamSize -
alignedTvList.getReservedMemoryBytes();
if (deltaBytes > 0) {
ownerQuery.getMemoryReservationContext().reserveMemoryCumulatively(deltaBytes);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java
index c0a71bf7edc..223e9ebb811 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/ReadOnlyMemChunk.java
@@ -136,11 +136,11 @@ public class ReadOnlyMemChunk {
int queryRowCount = entry.getValue();
if (!tvList.isSorted() && queryRowCount > tvList.seqRowCount()) {
entry.setValue(tvList.sort());
- long tvListRamSize = tvList.calculateRamSize().getRamSize();
tvList.lockQueryList();
try {
FragmentInstanceContext ownerQuery = (FragmentInstanceContext)
tvList.getOwnerQuery();
if (ownerQuery != null) {
+ long tvListRamSize = tvList.calculateRamSize().getRamSize();
long deltaBytes = tvListRamSize - tvList.getReservedMemoryBytes();
if (deltaBytes > 0) {
ownerQuery.getMemoryReservationContext().reserveMemoryCumulatively(deltaBytes);
@@ -288,11 +288,11 @@ public class ReadOnlyMemChunk {
int queryLength = entry.getValue();
if (!tvList.isSorted() && queryLength > tvList.seqRowCount()) {
entry.setValue(tvList.sort());
- long tvListRamSize = tvList.calculateRamSize().getRamSize();
tvList.lockQueryList();
try {
FragmentInstanceContext ownerQuery = (FragmentInstanceContext)
tvList.getOwnerQuery();
if (ownerQuery != null) {
+ long tvListRamSize = tvList.calculateRamSize().getRamSize();
long deltaBytes = tvListRamSize - tvList.getReservedMemoryBytes();
if (deltaBytes > 0) {
ownerQuery.getMemoryReservationContext().reserveMemoryCumulatively(deltaBytes);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
index f787ff64fda..d319b9c9a0d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/AlignedTVList.java
@@ -19,6 +19,8 @@
package org.apache.iotdb.db.utils.datastructure;
+import
org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext;
+import org.apache.iotdb.db.queryengine.execution.fragment.QueryContext;
import org.apache.iotdb.db.queryengine.plan.statement.component.Ordering;
import
org.apache.iotdb.db.storageengine.dataregion.wal.buffer.IWALByteBufferView;
import org.apache.iotdb.db.storageengine.dataregion.wal.utils.WALWriteUtils;
@@ -51,8 +53,10 @@ import java.io.DataInputStream;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.HashSet;
import java.util.List;
import java.util.Objects;
+import java.util.Set;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
@@ -78,6 +82,56 @@ public abstract class AlignedTVList extends TVList {
private long materializedBitmapMemoryCost;
private long arrayMemCostWithoutIndex;
+ /**
+ * A fully prepared partial clone. All allocations and validations are
completed before this plan
+ * is returned, so {@link #commit()} only moves already captured references
and updates primitive
+ * accounting fields.
+ */
+ public static final class PartialClonePlan {
+ private final AlignedTVList sourceList;
+ private final AlignedTVList cloneList;
+ private final List<Object>[] valueColumnsToMove;
+ private final List<BitMap>[] bitmapColumnsToMove;
+ private final long sourceArrayMemCostWithoutIndex;
+ private final long cloneArrayMemCostWithoutIndex;
+ private final long sourceBitmapMemoryCost;
+ private final long cloneBitmapMemoryCost;
+
+ private boolean committed;
+
+ private PartialClonePlan(
+ AlignedTVList sourceList,
+ AlignedTVList cloneList,
+ List<Object>[] valueColumnsToMove,
+ List<BitMap>[] bitmapColumnsToMove,
+ long sourceArrayMemCostWithoutIndex,
+ long cloneArrayMemCostWithoutIndex,
+ long sourceBitmapMemoryCost,
+ long cloneBitmapMemoryCost) {
+ this.sourceList = sourceList;
+ this.cloneList = cloneList;
+ this.valueColumnsToMove = valueColumnsToMove;
+ this.bitmapColumnsToMove = bitmapColumnsToMove;
+ this.sourceArrayMemCostWithoutIndex = sourceArrayMemCostWithoutIndex;
+ this.cloneArrayMemCostWithoutIndex = cloneArrayMemCostWithoutIndex;
+ this.sourceBitmapMemoryCost = sourceBitmapMemoryCost;
+ this.cloneBitmapMemoryCost = cloneBitmapMemoryCost;
+ }
+
+ public AlignedTVList getCloneList() {
+ return cloneList;
+ }
+
+ /** Commit the prepared ownership transfer. This method is idempotent and
allocation-free. */
+ public synchronized void commit() {
+ if (committed) {
+ return;
+ }
+ sourceList.commitPartialClone(this);
+ committed = true;
+ }
+ }
+
// Data type list -> list of TVList, add 1 when expanded -> primitive array
of basic type
// Index relation: columnIndex(dataTypeIndex) -> arrayIndex -> elementIndex
protected List<List<Object>> values;
@@ -96,12 +150,13 @@ public abstract class AlignedTVList extends TVList {
super();
dataTypes = types;
memoryBinaryChunkSize = new long[dataTypes.size()];
- refreshArrayMemCostWithoutIndex();
-
values = new ArrayList<>(types.size());
for (int i = 0; i < types.size(); i++) {
values.add(new ArrayList<>());
}
+ // arrayMemCostWithoutIndex depends on per-column value arrays, so values
must be
+ // initialized before computing it
+ refreshArrayMemCostWithoutIndex();
}
public static AlignedTVList newAlignedList(List<TSDataType> dataTypes) {
@@ -143,7 +198,7 @@ public abstract class AlignedTVList extends TVList {
alignedTvList.bitMaps = bitMaps;
alignedTvList.rowCount = this.rowCount;
alignedTvList.allValueColDeletedMap = getAllValueColDeletedMap();
- alignedTvList.materializedBitmapMemoryCost =
calculateBitmapRamCost(bitMaps);
+ alignedTvList.materializedBitmapMemoryCost =
calculateBitmapRamCost(bitMaps, null);
return alignedTvList;
}
@@ -162,34 +217,134 @@ public abstract class AlignedTVList extends TVList {
public synchronized AlignedTVList clone() {
AlignedTVList cloneList = AlignedTVList.newAlignedList(new
ArrayList<>(dataTypes));
cloneAs(cloneList);
- System.arraycopy(
- memoryBinaryChunkSize, 0, cloneList.memoryBinaryChunkSize, 0,
dataTypes.size());
- for (int i = 0; i < values.size(); i++) {
- // Clone value
+ cloneColumnDataTo(cloneList, null);
+ return cloneList;
+ }
+
+ /**
+ * Prepare a partial clone without changing this TVList. The returned plan
must be committed only
+ * after the query-memory reservation succeeds.
+ */
+ public synchronized PartialClonePlan preparePartialClone(Set<Integer>
columnsToClone) {
+ Set<Integer> retainedColumns =
+ new HashSet<>(Objects.requireNonNull(columnsToClone, "columnsToClone
cannot be null"));
+ AlignedTVList cloneList = AlignedTVList.newAlignedList(new
ArrayList<>(dataTypes));
+ cloneAs(cloneList);
+ cloneColumnDataTo(cloneList, retainedColumns);
+ return prepareMovePlan(cloneList, retainedColumns);
+ }
+
+ @SuppressWarnings("unchecked")
+ private PartialClonePlan prepareMovePlan(AlignedTVList cloneList,
Set<Integer> retainedColumns) {
+ Objects.requireNonNull(cloneList, "cloneList cannot be null");
+ int columnCount = values.size();
+ if (cloneList.values.size() != columnCount
+ || cloneList.memoryBinaryChunkSize.length !=
memoryBinaryChunkSize.length) {
+ throw new IllegalStateException("Target AlignedTVList has incompatible
column containers");
+ }
+
+ List<Object>[] valueColumnsToMove = (List<Object>[]) new
List<?>[columnCount];
+ List<BitMap>[] bitmapColumnsToMove = (List<BitMap>[]) new
List<?>[columnCount];
+ for (int i = 0; i < columnCount; i++) {
+ if (retainedColumns.contains(i)) {
+ continue;
+ }
+
List<Object> columnValues = values.get(i);
- for (Object valueArray : columnValues) {
- cloneList.values.get(i).add(cloneValue(dataTypes.get(i), valueArray));
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format("Missing value arrays for aligned column index %d
during move", i));
}
- // Clone bitmap in columnIndex
+ if (cloneList.values.get(i) == null ||
!cloneList.values.get(i).isEmpty()) {
+ throw new IllegalStateException(
+ String.format("Target value column index %d is not ready for
move", i));
+ }
+ valueColumnsToMove[i] = columnValues;
+
if (bitMaps != null && bitMaps.get(i) != null) {
- List<BitMap> columnBitMaps = bitMaps.get(i);
- if (cloneList.bitMaps == null) {
- cloneList.bitMaps = new ArrayList<>(dataTypes.size());
- for (int j = 0; j < dataTypes.size(); j++) {
- cloneList.bitMaps.add(null);
- }
- }
- if (cloneList.bitMaps.get(i) == null) {
- List<BitMap> cloneColumnBitMaps = new ArrayList<>();
- for (BitMap bitMap : columnBitMaps) {
- cloneColumnBitMaps.add(bitMap == null ? null : bitMap.clone());
- }
- cloneList.bitMaps.set(i, cloneColumnBitMaps);
+ if (cloneList.bitMaps == null
+ || cloneList.bitMaps.size() != bitMaps.size()
+ || cloneList.bitMaps.get(i) != null) {
+ throw new IllegalStateException(
+ String.format("Target bitmap column index %d is not ready for
move", i));
}
+ bitmapColumnsToMove[i] = bitMaps.get(i);
}
}
- cloneList.materializedBitmapMemoryCost = materializedBitmapMemoryCost;
- return cloneList;
+
+ return new PartialClonePlan(
+ this,
+ cloneList,
+ valueColumnsToMove,
+ bitmapColumnsToMove,
+ calculateArrayMemCostWithoutIndex(retainedColumns),
+ cloneList.calculateArrayMemCostWithoutIndex(null),
+ calculateBitmapRamCost(bitMaps, retainedColumns),
+ calculateBitmapRamCost(bitMaps, null));
+ }
+
+ private synchronized void commitPartialClone(PartialClonePlan plan) {
+ for (int i = 0; i < plan.valueColumnsToMove.length; i++) {
+ List<Object> columnValues = plan.valueColumnsToMove[i];
+ if (columnValues == null) {
+ continue;
+ }
+
+ plan.cloneList.values.set(i, columnValues);
+ values.set(i, null);
+ List<BitMap> columnBitMaps = plan.bitmapColumnsToMove[i];
+ if (columnBitMaps != null) {
+ plan.cloneList.bitMaps.set(i, columnBitMaps);
+ bitMaps.set(i, null);
+ }
+ memoryBinaryChunkSize[i] = 0;
+ }
+
+ arrayMemCostWithoutIndex = plan.sourceArrayMemCostWithoutIndex;
+ plan.cloneList.arrayMemCostWithoutIndex =
plan.cloneArrayMemCostWithoutIndex;
+ materializedBitmapMemoryCost = plan.sourceBitmapMemoryCost;
+ plan.cloneList.materializedBitmapMemoryCost = plan.cloneBitmapMemoryCost;
+ }
+
+ /**
+ * Release memory for non-query columns in this TVList. This is used during
memory ownership
+ * transfer from write process to read process to reduce memory footprint.
Only columns that are
+ * accessed by active queries are retained; all other columns are released.
+ *
+ * @param columnsToKeep set of column indices that are accessed by queries
and should be kept
+ */
+ public synchronized void releaseNonQueryColumns(Set<Integer> columnsToKeep) {
+ if (columnsToKeep == null || columnsToKeep.isEmpty()) {
+ return;
+ }
+
+ for (int i = 0; i < values.size(); i++) {
+ // Skip columns that should be kept or are already null
+ if (columnsToKeep.contains(i)) {
+ continue;
+ }
+
+ List<Object> columnValues = values.get(i);
+ if (columnValues == null) {
+ continue;
+ }
+
+ // Release memory for non-query columns
+ for (Object dataArray : columnValues) {
+ PrimitiveArrayManager.release(dataArray);
+ }
+ values.set(i, null);
+ memoryBinaryChunkSize[i] = 0;
+
+ // Release bitmap memory for non-query columns
+ if (bitMaps != null && bitMaps.get(i) != null) {
+ bitMaps.set(i, null);
+ }
+ }
+
+ materializedBitmapMemoryCost = calculateBitmapRamCost(bitMaps,
columnsToKeep);
+ // Refresh per-block memory cost after releasing columns
+ refreshArrayMemCostWithoutIndex();
}
@SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity
warning
@@ -203,10 +358,14 @@ public abstract class AlignedTVList extends TVList {
timestamps.get(arrayIndex)[elementIndex] = timestamp;
for (int i = 0; i < values.size(); i++) {
Object columnValue = value[i];
- List<Object> columnValues = values.get(i);
if (columnValue == null) {
markNullValue(i, arrayIndex, elementIndex);
}
+ List<Object> columnValues = values.get(i);
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format("Missing value arrays for aligned column index %d
during append", i));
+ }
switch (dataTypes.get(i)) {
case TEXT:
case BLOB:
@@ -547,6 +706,25 @@ public abstract class AlignedTVList extends TVList {
return dataTypes;
}
+ /**
+ * Get the union of all columns accessed by queries on this AlignedTVList.
This method should be
+ * called with queryListLock held for thread safety.
+ *
+ * @return set of accessed column indices, or empty set if no columns are
tracked or no queries
+ * are present
+ */
+ @Override
+ public Set<Integer> getAccessedColumnsForQuery() {
+ Set<Integer> accessedColumns = new HashSet<>();
+ for (QueryContext queryContext : getQueryContextSet()) {
+ if (queryContext instanceof FragmentInstanceContext) {
+ accessedColumns.addAll(
+ ((FragmentInstanceContext)
queryContext).getAccessedAlignedColumns(this));
+ }
+ }
+ return accessedColumns;
+ }
+
@Override
/*
* Must be synchronized with sort() on the same TVList instance: a query may
sort
@@ -604,6 +782,13 @@ public abstract class AlignedTVList extends TVList {
* or delete wrong rows.
*/
public synchronized void deleteColumn(int columnIndex) {
+ List<Object> columnValues = values.get(columnIndex);
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format(
+ "Missing value arrays for aligned column index %d during
delete", columnIndex));
+ }
+
if (bitMaps == null) {
List<List<BitMap>> localBitMaps = new ArrayList<>(dataTypes.size());
for (int j = 0; j < dataTypes.size(); j++) {
@@ -611,9 +796,10 @@ public abstract class AlignedTVList extends TVList {
}
bitMaps = localBitMaps;
}
+
if (bitMaps.get(columnIndex) == null) {
List<BitMap> columnBitMaps = new ArrayList<>();
- for (int i = 0; i < values.get(columnIndex).size(); i++) {
+ for (int i = 0; i < columnValues.size(); i++) {
columnBitMaps.add(new BitMap(ARRAY_SIZE));
}
bitMaps.set(columnIndex, columnBitMaps);
@@ -670,6 +856,69 @@ public abstract class AlignedTVList extends TVList {
}
}
+ /*
+ * There are two clone modes:
+ * 1. Full clone: columnsToClone is null, meaning no column filter is
applied. All columns are
+ * deep-cloned.
+ * 2. Partial clone: columnsToClone is non-null. Columns in columnsToClone
are deep-cloned for the
+ * query that keeps using the source TVList; columns not in
columnsToClone are not copied here.
+ * They are moved from the source TVList to cloneList later, and
cloneList becomes the new
+ * working list in the memtable.
+ *
+ * This method only performs the allocation phase: clone requested
value/bitmap arrays and prepare
+ * bitmap containers that will be needed by moved columns. It must not clear
or move columns from
+ * the source TVList here. The destructive move is performed only by
PartialClonePlan.commit()
+ * after cloneList and the ownership-transfer plan are fully prepared for
publication.
+ */
+ private void cloneColumnDataTo(AlignedTVList cloneList, Set<Integer>
columnsToClone) {
+ boolean cloneAllColumns = columnsToClone == null;
+ System.arraycopy(
+ memoryBinaryChunkSize, 0, cloneList.memoryBinaryChunkSize, 0,
dataTypes.size());
+ boolean hasBitMapsToMove = false;
+ for (int i = 0; i < values.size(); i++) {
+ // Clone value
+ List<Object> columnValues = values.get(i);
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format("Missing value arrays for aligned column index %d
during clone", i));
+ }
+ boolean shouldCloneColumn = cloneAllColumns ||
columnsToClone.contains(i);
+ if (!shouldCloneColumn) {
+ hasBitMapsToMove |= bitMaps != null && bitMaps.get(i) != null;
+ continue;
+ }
+
+ for (Object valueArray : columnValues) {
+ cloneList.values.get(i).add(cloneValue(dataTypes.get(i), valueArray));
+ }
+ // Clone bitmap in columnIndex
+ if (bitMaps != null && bitMaps.get(i) != null) {
+ List<BitMap> columnBitMaps = bitMaps.get(i);
+ if (cloneList.bitMaps == null) {
+ cloneList.bitMaps = new ArrayList<>(dataTypes.size());
+ for (int j = 0; j < dataTypes.size(); j++) {
+ cloneList.bitMaps.add(null);
+ }
+ }
+ if (cloneList.bitMaps.get(i) == null) {
+ List<BitMap> cloneColumnBitMaps = new ArrayList<>();
+ for (BitMap bitMap : columnBitMaps) {
+ cloneColumnBitMaps.add(bitMap == null ? null : bitMap.clone());
+ }
+ cloneList.bitMaps.set(i, cloneColumnBitMaps);
+ }
+ }
+ }
+ cloneList.materializedBitmapMemoryCost = materializedBitmapMemoryCost;
+
+ if (hasBitMapsToMove && cloneList.bitMaps == null) {
+ cloneList.bitMaps = new ArrayList<>(dataTypes.size());
+ for (int i = 0; i < dataTypes.size(); i++) {
+ cloneList.bitMaps.add(null);
+ }
+ }
+ }
+
@Override
protected void clearValue() {
for (int i = 0; i < dataTypes.size(); i++) {
@@ -703,7 +952,12 @@ public abstract class AlignedTVList extends TVList {
indices.add((int[]) getPrimitiveArraysByType(TSDataType.INT32));
}
for (int i = 0; i < dataTypes.size(); i++) {
- values.get(i).add(getPrimitiveArraysByType(dataTypes.get(i)));
+ List<Object> columnValues = values.get(i);
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format("Missing value arrays for aligned column index %d
during expand", i));
+ }
+ columnValues.add(getPrimitiveArraysByType(dataTypes.get(i)));
if (bitMaps != null && bitMaps.get(i) != null) {
bitMaps.get(i).add(null);
materializedBitmapMemoryCost += bitmapReferenceRamCost();
@@ -829,6 +1083,10 @@ public abstract class AlignedTVList extends TVList {
continue;
}
List<Object> columnValues = values.get(i);
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format("Missing value arrays for aligned column index %d
during arrayCopy", i));
+ }
switch (dataTypes.get(i)) {
case TEXT:
case BLOB:
@@ -871,6 +1129,14 @@ public abstract class AlignedTVList extends TVList {
}
private BitMap getBitMap(int columnIndex, int arrayIndex) {
+ List<Object> columnValues = values.get(columnIndex);
+ if (columnValues == null) {
+ throw new IllegalStateException(
+ String.format(
+ "Missing value arrays for aligned column index %d during mark
null value",
+ columnIndex));
+ }
+
// init BitMaps if doesn't have
if (bitMaps == null) {
List<List<BitMap>> localBitMaps = new ArrayList<>(dataTypes.size());
@@ -883,7 +1149,7 @@ public abstract class AlignedTVList extends TVList {
// if the bitmap in columnIndex is null, init the bitmap of this column
from the beginning
if (bitMaps.get(columnIndex) == null) {
List<BitMap> columnBitMaps = new ArrayList<>();
- for (int i = 0; i < values.get(columnIndex).size(); i++) {
+ for (int i = 0; i < columnValues.size(); i++) {
columnBitMaps.add(null);
}
bitMaps.set(columnIndex, columnBitMaps);
@@ -919,30 +1185,107 @@ public abstract class AlignedTVList extends TVList {
new ArrayList<>(dataTypes));
}
+ public synchronized RamInfo calculateRamSize(Set<Integer> columnsToClone) {
+ return new RamInfo(
+ timestamps.size(),
+ alignedTvListArrayMemCost(columnsToClone),
+ getRamSize(columnsToClone),
+ rowCount,
+ new ArrayList<>(dataTypes));
+ }
+
public synchronized long getRamSize() {
- return (long) timestamps.size()
+ return timestamps.size()
* (arrayMemCostWithoutIndex
+ (indices != null ? (long) PrimitiveArrayManager.ARRAY_SIZE *
Integer.BYTES : 0))
- + materializedBitmapMemoryCost;
+ + materializedBitmapMemoryCost
+ + calculateContainerRamCost(null);
+ }
+
+ public synchronized long getRamSize(Set<Integer> columnsToClone) {
+ return timestamps.size() * alignedTvListArrayMemCost(columnsToClone)
+ + calculateBitmapRamCost(bitMaps, columnsToClone)
+ + calculateContainerRamCost(columnsToClone);
+ }
+
+ private long calculateArrayMemCostWithoutIndex(Set<Integer> retainedColumns)
{
+ long arrayMemCost = alignedTvListArrayMemCost(retainedColumns);
+ if (indices != null) {
+ arrayMemCost -= (long) PrimitiveArrayManager.ARRAY_SIZE * Integer.BYTES;
+ }
+ return arrayMemCost;
}
private void refreshArrayMemCostWithoutIndex() {
- arrayMemCostWithoutIndex = alignedTvListArrayMemCost();
+ arrayMemCostWithoutIndex = calculateArrayMemCostWithoutIndex(null);
+ }
+
+ /**
+ * Calculate the one-time container memory retained by this list.
Primitive-array references in
+ * the time/index/value lists and bitmap lists are already charged by the
per-block accounting, so
+ * only their list objects and backing-array headers are added here. In
contrast, references in
+ * the outer column containers are not charged elsewhere and are counted in
full.
+ */
+ private long calculateContainerRamCost(Set<Integer> retainedColumns) {
+ long size = 0;
+
+ size += listRamCostWithReferences(dataTypes);
+ size += RamUsageEstimator.sizeOfLongArray(memoryBinaryChunkSize.length);
+ size += listRamCostWithoutReferences(timestamps);
if (indices != null) {
- arrayMemCostWithoutIndex -= (long) PrimitiveArrayManager.ARRAY_SIZE *
Integer.BYTES;
+ size += listRamCostWithoutReferences(indices);
+ }
+
+ size += listRamCostWithReferences(values);
+ for (int i = 0; i < values.size(); i++) {
+ if (retainedColumns != null && !retainedColumns.contains(i)) {
+ continue;
+ }
+ List<Object> columnValues = values.get(i);
+ if (columnValues != null) {
+ size += listRamCostWithoutReferences(columnValues);
+ }
+ }
+
+ if (bitMaps != null) {
+ size += listRamCostWithReferences(bitMaps);
+ for (int i = 0; i < bitMaps.size(); i++) {
+ if (retainedColumns != null && !retainedColumns.contains(i)) {
+ continue;
+ }
+ List<BitMap> columnBitMaps = bitMaps.get(i);
+ if (columnBitMaps != null) {
+ size += listRamCostWithoutReferences(columnBitMaps);
+ }
+ }
}
+ return size;
}
- private static long calculateBitmapRamCost(List<List<BitMap>> bitMaps) {
+ private static long listRamCostWithReferences(List<?> list) {
+ return RamUsageEstimator.shallowSizeOf(list) +
RamUsageEstimator.sizeOfObjectArray(list.size());
+ }
+
+ private static long listRamCostWithoutReferences(List<?> list) {
+ return RamUsageEstimator.shallowSizeOf(list)
+ + (list.isEmpty() ? 0 : RamUsageEstimator.sizeOfObjectArray(0));
+ }
+
+ private static long calculateBitmapRamCost(
+ List<List<BitMap>> bitMaps, Set<Integer> columnsToClone) {
if (bitMaps == null) {
return 0;
}
long size = 0;
- for (List<BitMap> columnBitMaps : bitMaps) {
+ for (int i = 0, length = bitMaps.size(); i < length; i++) {
+ if (columnsToClone != null && !columnsToClone.contains(i)) {
+ continue;
+ }
+ List<BitMap> columnBitMaps = bitMaps.get(i);
if (columnBitMaps == null) {
continue;
}
- size += (long) columnBitMaps.size() * bitmapReferenceRamCost();
+ size += columnBitMaps.size() * bitmapReferenceRamCost();
for (BitMap bitMap : columnBitMaps) {
if (bitMap != null) {
size += bitmapRamCost();
@@ -986,30 +1329,36 @@ public abstract class AlignedTVList extends TVList {
*
* @return AlignedTvListArrayMemSize
*/
- public long alignedTvListArrayMemCost() {
+ public long alignedTvListArrayMemCost(Set<Integer> columnsToClone) {
long size = 0;
+ int retainedColumnNum = 0;
// value array mem size
for (int column = 0; column < dataTypes.size(); column++) {
+ if (columnsToClone != null && !columnsToClone.contains(column)) {
+ continue;
+ }
TSDataType type = dataTypes.get(column);
- if (type != null) {
+ if (type != null && values.get(column) != null) {
+ retainedColumnNum++;
size += (long) PrimitiveArrayManager.ARRAY_SIZE * (long)
type.getDataTypeSize();
}
}
- // size is 0 when all types are null
- if (size == 0) {
- return size;
- }
+
// time array mem size
size += PrimitiveArrayManager.ARRAY_SIZE * 8L;
// index array mem size
size += (indices != null) ? PrimitiveArrayManager.ARRAY_SIZE * 4L : 0;
// array headers mem size
- size += (long) NUM_BYTES_ARRAY_HEADER * (2 + dataTypes.size());
+ size += (long) NUM_BYTES_ARRAY_HEADER * (2 + retainedColumnNum);
// Object references size in ArrayList
- size += (long) NUM_BYTES_OBJECT_REF * (2 + dataTypes.size());
+ size += (long) NUM_BYTES_OBJECT_REF * (2 + retainedColumnNum);
return size;
}
+ public long alignedTvListArrayMemCost() {
+ return alignedTvListArrayMemCost((Set<Integer>) null);
+ }
+
/**
* Get the single column array mem cost by give type.
*
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
index a085c29e130..60e421c9af5 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/datastructure/TVList.java
@@ -809,6 +809,16 @@ public abstract class TVList implements WALEntryValue {
return queryContextSet;
}
+ /**
+ * Get the union of all columns accessed by queries on this TVList. For
non-AlignedTVList, returns
+ * empty set. This method should be called with queryListLock held for
thread safety.
+ *
+ * @return set of accessed column indices, or empty set if no columns are
tracked
+ */
+ public Set<Integer> getAccessedColumnsForQuery() {
+ return null;
+ }
+
public List<BitMap> getBitMap() {
return bitMap;
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceExecutionTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceExecutionTest.java
index cfc7f887dcf..0f1b1c7d253 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceExecutionTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceExecutionTest.java
@@ -49,6 +49,7 @@ import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.file.metadata.enums.CompressionType;
import org.apache.tsfile.file.metadata.enums.TSEncoding;
import org.apache.tsfile.read.reader.IPointReader;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.junit.Test;
import org.mockito.Mockito;
@@ -298,4 +299,24 @@ public class FragmentInstanceExecutionTest {
}
return memTable;
}
+
+ private IMemTable createMemTable(String deviceId, List<IMeasurementSchema>
schemaList)
+ throws IllegalPathException {
+ PrimitiveMemTable memTable = new PrimitiveMemTable("root.test", "1");
+
+ // Insert data in reverse order to make it unsorted
+ int rows = 100;
+ for (int i = rows - 1; i >= 0; i--) {
+ Object[] values = new Object[5];
+ for (int j = 0; j < 5; j++) {
+ values[j] = (long) i * 100 + j;
+ }
+ memTable.writeAlignedRow(
+ DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)),
+ schemaList,
+ i,
+ values);
+ }
+ return memTable;
+ }
}
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
index 6d0cabb0443..9768d1177d5 100644
---
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
@@ -20,6 +20,7 @@
package org.apache.iotdb.db.queryengine.plan.planner;
import org.apache.iotdb.db.queryengine.common.QueryId;
+import org.apache.iotdb.db.queryengine.exception.MemoryNotEnoughException;
import
org.apache.iotdb.db.queryengine.plan.planner.memory.NotThreadSafeMemoryReservationManager;
import org.junit.After;
@@ -161,4 +162,49 @@ public class LocalExecutionPlannerOperatorsMemoryTest {
Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
}
+
+ @Test
+ public void testImmediateReservationRollback() {
+ 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.reserveMemoryCumulatively(request);
+ manager.releaseMemoryImmediately(request);
+ manager.reserveMemoryImmediately();
+
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+
+ manager.reserveMemoryImmediately(request);
+ manager.releaseMemoryImmediately(request);
+
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
+
+ @Test
+ public void testFailedCumulativeReservationDoesNotRemainPending() {
+ long freeBefore = PLANNER.getFreeMemoryForOperators();
+ long request = freeBefore + MEMORY_BATCH_THRESHOLD;
+ NotThreadSafeMemoryReservationManager manager =
+ new NotThreadSafeMemoryReservationManager(new QueryId("normal_query"),
"test");
+
+ try {
+ manager.reserveMemoryCumulatively(request);
+ Assert.fail("Expected insufficient query memory");
+ } catch (MemoryNotEnoughException expected) {
+ // expected
+ }
+
+ // A stale pending reservation would make this retry fail again.
+ manager.reserveMemoryImmediately();
+ Assert.assertEquals(0L, manager.getReservedBytesInTotalForTest());
+ Assert.assertEquals(freeBefore, PLANNER.getFreeMemoryForOperators());
+ }
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtilsTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtilsTest.java
new file mode 100644
index 00000000000..cd1ece64f08
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtilsTest.java
@@ -0,0 +1,109 @@
+/*
+ * 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.schemaengine.schemaregion.utils;
+
+import org.apache.iotdb.commons.path.MeasurementPath;
+import org.apache.iotdb.db.queryengine.execution.fragment.QueryContext;
+import org.apache.iotdb.db.storageengine.dataregion.memtable.IWritableMemChunk;
+import org.apache.iotdb.db.utils.datastructure.TVList;
+
+import org.apache.tsfile.enums.TSDataType;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public class ResourceByPathUtilsTest {
+
+ @Test
+ public void testFlushingQueryLocksTemporaryTVListBeforeRegistration() throws
Exception {
+ TVList candidate = TVList.newList(TSDataType.INT64);
+ candidate.putLong(2, 2);
+ candidate.putLong(1, 1);
+ Assert.assertFalse(candidate.isSorted());
+
+ QueryContext previousQuery = new QueryContext(1, false);
+ candidate.lockQueryList();
+ try {
+ candidate.getQueryContextSet().add(previousQuery);
+ } finally {
+ candidate.unlockQueryList();
+ }
+
+ TVList temporaryList = candidate.cloneForFlushSort();
+ IWritableMemChunk memChunk = mock(IWritableMemChunk.class);
+ when(memChunk.getSortedList()).thenReturn(Collections.emptyList());
+ when(memChunk.getWorkingTVList()).thenReturn(candidate);
+ CountDownLatch temporaryListInitialized = new CountDownLatch(1);
+ when(memChunk.initWorkingListForFlushIfNecessary(candidate, true))
+ .thenAnswer(
+ ignored -> {
+ temporaryListInitialized.countDown();
+ return temporaryList;
+ });
+
+ ResourceByPathUtils resourceByPathUtils =
+ ResourceByPathUtils.getResourceInstance(
+ new MeasurementPath("root.test.d.s", TSDataType.INT64));
+ QueryContext currentQuery = new QueryContext(2, false);
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ Future<Map<TVList, Integer>> result = null;
+ try {
+ temporaryList.lockQueryList();
+ try {
+ result =
+ executor.submit(
+ () ->
+ resourceByPathUtils.prepareTvListMapForQuery(
+ currentQuery, memChunk, false, null, null));
+ Assert.assertTrue(temporaryListInitialized.await(3, TimeUnit.SECONDS));
+ Future<Map<TVList, Integer>> blockedResult = result;
+ Assert.assertThrows(
+ TimeoutException.class, () -> blockedResult.get(200,
TimeUnit.MILLISECONDS));
+ } finally {
+ temporaryList.unlockQueryList();
+ }
+
+ Map<TVList, Integer> tvListQueryMap = result.get(3, TimeUnit.SECONDS);
+ Assert.assertTrue(tvListQueryMap.containsKey(temporaryList));
+ temporaryList.lockQueryList();
+ try {
+
Assert.assertTrue(temporaryList.getQueryContextSet().contains(currentQuery));
+ } finally {
+ temporaryList.unlockQueryList();
+ }
+ } finally {
+ if (result != null) {
+ result.cancel(true);
+ }
+ executor.shutdownNow();
+ executor.awaitTermination(3, TimeUnit.SECONDS);
+ }
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
index db3b84dbf58..0e2e533cc87 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java
@@ -688,11 +688,11 @@ public class TsFileProcessorTest {
// Test Tablet
processor.insertTablet(genInsertTableNode(0, true), 0, 10, new
TSStatus[10]);
IMemTable memTable = processor.getWorkMemTable();
- Assert.assertEquals(1596552, memTable.getTVListsRamCost());
+ Assert.assertEquals(1764688, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNode(100, true), 0, 10, new
TSStatus[10]);
- Assert.assertEquals(1596552, memTable.getTVListsRamCost());
+ Assert.assertEquals(1764688, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNode(200, true), 0, 10, new
TSStatus[10]);
- Assert.assertEquals(1596552, memTable.getTVListsRamCost());
+ Assert.assertEquals(1764688, memTable.getTVListsRamCost());
Assert.assertEquals(90000, memTable.getTotalPointsNum());
Assert.assertEquals(720360, memTable.memSize());
// Test records
@@ -701,7 +701,7 @@ public class TsFileProcessorTest {
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
- Assert.assertEquals(1598168, memTable.getTVListsRamCost());
+ Assert.assertEquals(1766304, memTable.getTVListsRamCost());
Assert.assertEquals(90100, memTable.getTotalPointsNum());
Assert.assertEquals(721560, memTable.memSize());
}
@@ -724,21 +724,21 @@ public class TsFileProcessorTest {
// Test Tablet
processor.insertTablet(genInsertTableNode(0, true), 0, 10, new
TSStatus[10]);
IMemTable memTable = processor.getWorkMemTable();
- Assert.assertEquals(1596552, memTable.getTVListsRamCost());
+ Assert.assertEquals(1764688, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNodeFors3000ToS6000(0, true), 0, 10,
new TSStatus[10]);
- Assert.assertEquals(3552552, memTable.getTVListsRamCost());
+ Assert.assertEquals(4152728, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNode(100, true), 0, 10, new
TSStatus[10]);
- Assert.assertEquals(3552552, memTable.getTVListsRamCost());
+ Assert.assertEquals(4152728, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNodeFors3000ToS6000(100, true), 0,
10, new TSStatus[10]);
- Assert.assertEquals(3552552, memTable.getTVListsRamCost());
+ Assert.assertEquals(4152728, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNode(200, true), 0, 10, new
TSStatus[10]);
- Assert.assertEquals(3552552, memTable.getTVListsRamCost());
+ Assert.assertEquals(4152728, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNodeFors3000ToS6000(200, true), 0,
10, new TSStatus[10]);
- Assert.assertEquals(3552552, memTable.getTVListsRamCost());
+ Assert.assertEquals(4152728, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNode(300, true), 0, 10, new
TSStatus[10]);
- Assert.assertEquals(6937104, memTable.getTVListsRamCost());
+ Assert.assertEquals(7537280, memTable.getTVListsRamCost());
processor.insertTablet(genInsertTableNodeFors3000ToS6000(300, true), 0,
10, new TSStatus[10]);
- Assert.assertEquals(7105104, memTable.getTVListsRamCost());
+ Assert.assertEquals(7705280, memTable.getTVListsRamCost());
Assert.assertEquals(240000, memTable.getTotalPointsNum());
Assert.assertEquals(1920960, memTable.memSize());
@@ -748,14 +748,14 @@ public class TsFileProcessorTest {
record.addTuple(DataPoint.getDataPoint(dataType, measurementId,
String.valueOf(i)));
processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
- Assert.assertEquals(7106720, memTable.getTVListsRamCost());
+ Assert.assertEquals(7706896, memTable.getTVListsRamCost());
// Test records
for (int i = 1; i <= 100; i++) {
TSRecord record = new TSRecord(i, deviceId);
record.addTuple(DataPoint.getDataPoint(dataType, "s1",
String.valueOf(i)));
processor.insert(buildInsertRowNodeByTSRecord(record), new long[4]);
}
- Assert.assertEquals(7108336, memTable.getTVListsRamCost());
+ Assert.assertEquals(7708512, memTable.getTVListsRamCost());
Assert.assertEquals(240200, memTable.getTotalPointsNum());
Assert.assertEquals(1923360, memTable.memSize());
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java
index e27dc009f63..f84710bf6a0 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java
@@ -29,7 +29,10 @@ import org.junit.Test;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
import java.util.List;
+import java.util.Set;
import static
org.apache.iotdb.db.storageengine.rescon.memory.PrimitiveArrayManager.ARRAY_SIZE;
import static org.apache.tsfile.utils.RamUsageEstimator.NUM_BYTES_ARRAY_HEADER;
@@ -181,11 +184,11 @@ public class AlignedTVListTest {
ARRAY_SIZE / Byte.SIZE + 1,
firstColumnBitMaps.get(2).getByteArray().length);
Assert.assertTrue(tvList.isNullValue(ARRAY_SIZE * 2 + 1, 0));
Assert.assertFalse(tvList.isNullValue(ARRAY_SIZE * 2, 0));
- Assert.assertEquals(
+ long primitiveArrayAndBitmapCost =
3L * tvList.alignedTvListArrayMemCost()
+ 3L * AlignedTVList.bitmapReferenceRamCost()
- + AlignedTVList.bitmapRamCost(),
- tvList.getRamSize());
+ + AlignedTVList.bitmapRamCost();
+ Assert.assertTrue(tvList.getRamSize() > primitiveArrayAndBitmapCost);
Assert.assertEquals(tvList.getRamSize(),
tvList.calculateRamSize().getRamSize());
Assert.assertEquals(tvList.getRamSize(), tvList.clone().getRamSize());
Assert.assertEquals(tvList.getRamSize(),
tvList.cloneForFlushSort().getRamSize());
@@ -202,14 +205,14 @@ public class AlignedTVListTest {
long ramSizeBeforeExtension = tvList.getRamSize();
tvList.extendColumn(TSDataType.INT32);
- Assert.assertEquals(
- 2L
- * (AlignedTVList.valueListArrayMemCost(TSDataType.INT32)
- + AlignedTVList.bitmapReferenceRamCost()
- + AlignedTVList.bitmapRamCost()),
- tvList.getRamSize() - ramSizeBeforeExtension);
+ Assert.assertTrue(
+ tvList.getRamSize() - ramSizeBeforeExtension
+ >= 2L
+ * (AlignedTVList.valueListArrayMemCost(TSDataType.INT32)
+ + AlignedTVList.bitmapReferenceRamCost()
+ + AlignedTVList.bitmapRamCost()));
tvList.clear();
- Assert.assertEquals(0, tvList.getRamSize());
+ Assert.assertTrue(tvList.getRamSize() > 0);
}
@Test
@@ -334,4 +337,162 @@ public class AlignedTVListTest {
Assert.assertEquals(tvList.memoryBinaryChunkSize[1], 0);
Assert.assertEquals(tvList.memoryBinaryChunkSize[2], 0);
}
+
+ @Test
+ public void testMovesUnclonedColumns() {
+ List<TSDataType> dataTypes = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ dataTypes.add(TSDataType.INT64);
+ }
+ AlignedTVList tvList = AlignedTVList.newAlignedList(dataTypes);
+ tvList.putAlignedValue(0, new Object[] {1L, 2L, null});
+
+ Set<Integer> columnsToClone = Collections.singleton(1);
+ long retainedRamSize =
tvList.calculateRamSize(columnsToClone).getRamSize();
+ AlignedTVList.PartialClonePlan partialClonePlan =
tvList.preparePartialClone(columnsToClone);
+ AlignedTVList clonedTvList = partialClonePlan.getCloneList();
+
+ Assert.assertNotNull(tvList.getValues().get(0));
+ Assert.assertNotNull(tvList.getValues().get(2));
+ Assert.assertEquals(1L, tvList.getLongByValueIndex(0, 0));
+ Assert.assertTrue(tvList.isNullValue(0, 2));
+ Assert.assertEquals(2L, clonedTvList.getLongByValueIndex(0, 1));
+
+ partialClonePlan.commit();
+
+ Assert.assertNull(tvList.getValues().get(0));
+ Assert.assertNull(tvList.getValues().get(2));
+ Assert.assertTrue(tvList.isNullValue(0, 0));
+ Assert.assertTrue(tvList.isNullValue(0, 2));
+ Assert.assertEquals(1L, clonedTvList.getLongByValueIndex(0, 0));
+ Assert.assertEquals(2L, clonedTvList.getLongByValueIndex(0, 1));
+ Assert.assertTrue(clonedTvList.isNullValue(0, 2));
+ Assert.assertEquals(retainedRamSize,
tvList.calculateRamSize().getRamSize());
+ }
+
+ @Test
+ public void testPartialRamSizeIncludesWideColumnContainers() {
+ int columnCount = 256;
+ List<TSDataType> dataTypes = new ArrayList<>(columnCount);
+ Object[] values = new Object[columnCount];
+ for (int i = 0; i < columnCount; i++) {
+ dataTypes.add(TSDataType.INT64);
+ values[i] = (long) i;
+ }
+
+ AlignedTVList tvList = AlignedTVList.newAlignedList(dataTypes);
+ tvList.putAlignedValue(1, values);
+ Set<Integer> retainedColumns = Collections.singleton(0);
+ long primitiveArrayCost =
+ (long) tvList.getTimestamps().size() *
tvList.alignedTvListArrayMemCost(retainedColumns);
+ long retainedRamSize =
tvList.calculateRamSize(retainedColumns).getRamSize();
+
+ // memoryBinaryChunkSize and the outer column containers remain N-wide
after partial move.
+ Assert.assertTrue(retainedRamSize - primitiveArrayCost >= (long)
columnCount * Long.BYTES);
+
+ AlignedTVList.PartialClonePlan plan =
tvList.preparePartialClone(retainedColumns);
+ plan.commit();
+ Assert.assertEquals(retainedRamSize,
tvList.calculateRamSize().getRamSize());
+ }
+
+ @Test
+ public void testPartialReservationMatchesCleanupCalculation() {
+ for (boolean createIndices : new boolean[] {false, true}) {
+ for (boolean retainValueColumn : new boolean[] {false, true}) {
+ AlignedTVList tvList =
+ AlignedTVList.newAlignedList(
+ new ArrayList<>(
+ Arrays.asList(TSDataType.INT64, TSDataType.INT64,
TSDataType.INT64)));
+ for (int i = 0; i <= ARRAY_SIZE; i++) {
+ long time = createIndices ? ARRAY_SIZE - i : i;
+ tvList.putAlignedValue(
+ time, new Object[] {(long) i, i % 2 == 0 ? null : (long) i,
(long) i});
+ }
+ if (createIndices) {
+ Assert.assertFalse(tvList.isSorted());
+ tvList.sort();
+ Assert.assertNotNull(tvList.getIndices());
+ } else {
+ Assert.assertNull(tvList.getIndices());
+ }
+
+ Set<Integer> retainedColumns =
+ retainValueColumn ? Collections.singleton(1) :
Collections.emptySet();
+ long reservedMemoryBytes =
tvList.calculateRamSize(retainedColumns).getRamSize();
+ tvList.setReservedMemoryBytes(reservedMemoryBytes);
+
+ AlignedTVList.PartialClonePlan plan =
tvList.preparePartialClone(retainedColumns);
+ plan.commit();
+
+ long cleanupMemoryBytes = tvList.calculateRamSize().getRamSize();
+ String scenario =
+ String.format(
+ "createIndices=%s, retainValueColumn=%s", createIndices,
retainValueColumn);
+ Assert.assertEquals(scenario, reservedMemoryBytes, cleanupMemoryBytes);
+ Assert.assertEquals(scenario, tvList.getReservedMemoryBytes(),
cleanupMemoryBytes);
+ }
+ }
+ }
+
+ @Test
+ public void testPartialCloneFailureLeavesSourceUntouched() {
+ AlignedTVList tvList =
+ AlignedTVList.newAlignedList(
+ Arrays.asList(TSDataType.INT64, TSDataType.INT64,
TSDataType.INT64));
+ tvList.putAlignedValue(0, new Object[] {null, 2L, 3L});
+
+ List<Object> firstColumnValues = tvList.getValues().get(0);
+ List<Object> secondColumnValues = tvList.getValues().get(1);
+ List<Object> thirdColumnValues = tvList.getValues().get(2);
+ List<BitMap> firstColumnBitMaps = tvList.getBitMaps().get(0);
+ Object invalidThirdColumnArray = new int[ARRAY_SIZE];
+ thirdColumnValues.set(0, invalidThirdColumnArray);
+
+ Set<Integer> columnsToClone = new HashSet<>(Arrays.asList(0, 1, 2));
+ Assert.assertThrows(ClassCastException.class, () ->
tvList.preparePartialClone(columnsToClone));
+
+ Assert.assertSame(firstColumnValues, tvList.getValues().get(0));
+ Assert.assertSame(secondColumnValues, tvList.getValues().get(1));
+ Assert.assertSame(thirdColumnValues, tvList.getValues().get(2));
+ Assert.assertSame(invalidThirdColumnArray,
tvList.getValues().get(2).get(0));
+ Assert.assertSame(firstColumnBitMaps, tvList.getBitMaps().get(0));
+ Assert.assertTrue(tvList.isNullValue(0, 0));
+ Assert.assertEquals(2L, tvList.getLongByValueIndex(0, 1));
+ }
+
+ @Test
+ public void testReleaseNonQueryColumnsWithBitmaps() {
+ List<TSDataType> dataTypes = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ dataTypes.add(TSDataType.INT64);
+ }
+ AlignedTVList tvList = AlignedTVList.newAlignedList(dataTypes);
+ for (int i = 0; i < 100; i++) {
+ Object[] values = new Object[3];
+ values[0] = (long) i;
+ values[1] = null; // This will create a bitmap
+ values[2] = (long) (i * 100);
+ tvList.putAlignedValue(i, values);
+ }
+
+ // Verify bitmaps were created for column 1
+ Assert.assertNotNull(tvList.getBitMaps());
+ Assert.assertNotNull(tvList.getBitMaps().get(1));
+
+ // Keep only column 0 and 2, release column 1
+ Set<Integer> columnsToKeep = new HashSet<>(Arrays.asList(0, 2));
+ tvList.releaseNonQueryColumns(columnsToKeep);
+
+ // Verify column 1 is released
+ Assert.assertNull(tvList.getValues().get(1));
+ Assert.assertNull(tvList.getBitMaps().get(1));
+
+ // Verify columns 0 and 2 are intact
+ Assert.assertFalse(tvList.getValues().get(0).isEmpty());
+ Assert.assertFalse(tvList.getValues().get(2).isEmpty());
+ for (int i = 0; i < 100; i++) {
+ Assert.assertEquals((long) i, tvList.getLongByValueIndex(i, 0));
+ Assert.assertEquals((long) (i * 100), tvList.getLongByValueIndex(i, 2));
+ }
+ }
}