From edeb7c5276140dd5fc52f911a46dd3476cfb0fc5 Mon Sep 17 00:00:00 2001 From: shuwenwei <55970239+shuwenwei@users.noreply.github.com> Date: Wed, 5 Aug 2026 20:46:06 +0800 Subject: [PATCH 1/6] [To master] clone partial columns of aligned tvlist for query (cherry-picked from #18394) --- .../memory/MemoryReservationManager.java | 6 + .../fragment/FragmentInstanceContext.java | 50 +- .../memory/FakedMemoryReservationManager.java | 3 + ...NotThreadSafeMemoryReservationManager.java | 16 +- .../ThreadSafeMemoryReservationManager.java | 5 + .../utils/ResourceByPathUtils.java | 339 ++++++++++---- .../memtable/AbstractWritableMemChunk.java | 24 +- .../memtable/AlignedReadOnlyMemChunk.java | 4 +- .../dataregion/memtable/ReadOnlyMemChunk.java | 4 +- .../db/utils/datastructure/AlignedTVList.java | 443 ++++++++++++++++-- .../iotdb/db/utils/datastructure/TVList.java | 10 + .../FragmentInstanceExecutionTest.java | 21 + ...alExecutionPlannerOperatorsMemoryTest.java | 46 ++ .../utils/ResourceByPathUtilsTest.java | 109 +++++ .../datastructure/AlignedTVListTest.java | 161 +++++++ 15 files changed, 1082 insertions(+), 159 deletions(-) create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtilsTest.java diff --git a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/memory/MemoryReservationManager.java b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/memory/MemoryReservationManager.java index f0420330652e8..18ad895bb3146 100644 --- a/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/memory/MemoryReservationManager.java +++ b/iotdb-core/calc-commons/src/main/java/org/apache/iotdb/calc/plan/planner/memory/MemoryReservationManager.java @@ -48,6 +48,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/execution/fragment/FragmentInstanceContext.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/fragment/FragmentInstanceContext.java index e29e2a2141aa1..20ee72506c91e 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 @@ -78,8 +78,10 @@ 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.LinkedBlockingQueue; import java.util.concurrent.atomic.AtomicLong; @@ -184,6 +186,9 @@ public class FragmentInstanceContext extends QueryContext { private long closedUnseqFileNum = 0; private boolean highestPriority = false; + // accessed value columns on each referenced AlignedTVList. + private final Map> alignedTVListColumnAccessMap = new ConcurrentHashMap<>(); + public static FragmentInstanceContext createFragmentInstanceContext( FragmentInstanceId id, FragmentInstanceStateMachine stateMachine, @@ -253,6 +258,48 @@ public boolean isExternalTsFileScan() { return queryDataSourceType == QueryDataSourceType.EXTERNAL_TSFILE_SCAN; } + /** + * 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 columnIndexList) { + Set 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 getAccessedAlignedColumns(TVList tvList) { + Set accessedColumns = alignedTVListColumnAccessMap.get(tvList); + return accessedColumns == null + ? Collections.emptySet() + : Collections.unmodifiableSet(accessedColumns); + } + @TestOnly public static FragmentInstanceContext createFragmentInstanceContext( FragmentInstanceId id, FragmentInstanceStateMachine stateMachine) { @@ -1041,12 +1088,12 @@ public void releaseResourceWhenAllDriversAreClosed() { */ private void releaseTVListOwnedByQuery() { for (TVList tvList : tvListSet) { - long tvListRamSize = tvList.calculateRamSize().getRamSize(); tvList.lockQueryList(); Set queryContextSet = tvList.getQueryContextSet(); try { queryContextSet.remove(this); if (tvList.getOwnerQuery() == this) { + long tvListRamSize = tvList.calculateRamSize().getRamSize(); if (tvList.getReservedMemoryBytes() != tvListRamSize) { LOGGER.warn( DataNodeQueryMessages @@ -1131,6 +1178,7 @@ public synchronized void releaseResource() { // release TVList/AlignedTVList owned by current query releaseTVListOwnedByQuery(); + alignedTVListColumnAccessMap.clear(); fileModCache = null; tables = 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 d1c34e365efbc..8b45582f138a9 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 @@ -37,6 +37,9 @@ public void reserveMemoryImmediately(final long size) {} @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/NotThreadSafeMemoryReservationManager.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/memory/NotThreadSafeMemoryReservationManager.java index 71924894c7c96..39e987a4b2f8d 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 @@ -80,7 +80,14 @@ public long getFallbackBytesInTotalForTest() { 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; + } } } @@ -129,6 +136,13 @@ public void releaseMemoryCumulatively(final long size) { } } + @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 0a1c6eee4181e..71676e5b77fe9 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 synchronized void releaseMemoryCumulatively(long size) { 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 17d5dfd2995c7..8c7793896c736 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 @@ -40,6 +40,7 @@ import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; import org.apache.iotdb.db.utils.ModificationUtils; import org.apache.iotdb.db.utils.SchemaUtils; +import org.apache.iotdb.db.utils.datastructure.AlignedTVList; import org.apache.iotdb.db.utils.datastructure.TVList; import org.apache.tsfile.enums.TSDataType; @@ -74,6 +75,7 @@ import java.util.List; import java.util.Map; import java.util.stream.Collectors; +import java.util.Set; import static org.apache.iotdb.commons.path.AlignedPath.VECTOR_PLACEHOLDER; @@ -129,7 +131,8 @@ protected Map prepareTvListMapForQuery( QueryContext context, IWritableMemChunk memChunk, boolean isWorkMemTable, - Filter globalTimeFilter) { + Filter globalTimeFilter, + List columnIndexList) { // should copy globalTimeFilter because GroupByMonthFilter is stateful Filter copyTimeFilter = null; if (globalTimeFilter != null) { @@ -150,121 +153,249 @@ protected Map prepareTvListMapForQuery( .SCHEMA_LOG_FLUSHING_WORKING_MEMTABLE_ADD_CURRENT_QUERY_CONTEXT_TO_IMMUTABLE_7B7CD373); 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( - DataNodeSchemaMessages - .SCHEMA_LOG_FLUSHING_MEMTABLE_ADD_CURRENT_QUERY_CONTEXT_TO_MUTABLE_TVLIST_BEB0D766); - 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()); - } - } else { - if (list.isSorted() || list.getQueryContextSet().isEmpty()) { + boolean canUseListDirectly = + candidate.isSorted() || candidate.getQueryContextSet().isEmpty(); LOGGER.debug( DataNodeSchemaMessages - .SCHEMA_LOG_WORKING_MEMTABLE_ADD_CURRENT_QUERY_CONTEXT_TO_MUTABLE_TVLIST_8C937414); - list.getQueryContextSet().add(context); - tvListQueryMap.put(list, list.rowCount()); - } else { - /* - * +----------------------+ - * | MemTable | - * | | - * | +------------+ | +-----------------+ - * | | TVList |<---+--+ +---+ Previous Query | - * | +-----^------+ | | | +-----------------+ - * | | | | | - * +----------+-----------+ | | +----------------+ - * | Clone +---+---+ Current Query | - * +-----+------+ | +----------------+ - * | TVList | <---------+ - * +------------+ - */ + .SCHEMA_LOG_FLUSHING_MEMTABLE_ADD_CURRENT_QUERY_CONTEXT_TO_MUTABLE_TVLIST_BEB0D766); + 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; + } + + if (candidate.isSorted() || candidate.getQueryContextSet().isEmpty()) { LOGGER.debug( DataNodeSchemaMessages - .SCHEMA_LOG_WORKING_MEMTABLE_CLONE_MUTABLE_TVLIST_AND_REPLACE_OLD_TVLIST_FD1EAE22); - 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()); + .SCHEMA_LOG_WORKING_MEMTABLE_ADD_CURRENT_QUERY_CONTEXT_TO_MUTABLE_TVLIST_8C937414); + 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( + DataNodeSchemaMessages + .SCHEMA_LOG_WORKING_MEMTABLE_CLONE_MUTABLE_TVLIST_AND_REPLACE_OLD_TVLIST_FD1EAE22); + + 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 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( + DataNodeSchemaMessages.FAILED_TO_RESERVE_MEMORY_TVLIST, + listRamInfo.getRamSize(), + listRamInfo.getTimestampsSize(), + listRamInfo.getArrayMemCost(), + listRamInfo.getRowCount(), + listRamInfo.getDataTypes()); + } + throw ex; + } finally { + candidate.unlockQueryList(); } - } catch (MemoryNotEnoughException ex) { - LOGGER.warn( - DataNodeSchemaMessages.FAILED_TO_RESERVE_MEMORY_TVLIST, - listRamInfo.getRamSize(), - listRamInfo.getTimestampsSize(), - listRamInfo.getArrayMemCost(), - listRamInfo.getRowCount(), - listRamInfo.getDataTypes()); - throw ex; - } finally { - list.unlockQueryList(); } - if (cloneList != null) { - memChunk.setWorkingTVList(cloneList); - } - return tvListQueryMap; } } @@ -451,11 +582,6 @@ public ReadOnlyMemChunk getReadOnlyMemChunkFromMemTable( } } - // prepare AlignedTVList for query. It should clone TVList if necessary. - Map alignedTvListQueryMap = - prepareTvListMapForQuery( - context, alignedMemChunk, modsToMemtable == null, globalTimeFilter); - // column index list for the query // Columns with inconsistent types will be ignored and set -1 List columnIndexList = @@ -463,6 +589,11 @@ public ReadOnlyMemChunk getReadOnlyMemChunkFromMemTable( List timeColumnDeletion = null; List> valueColumnsDeletionList = null; + // prepare AlignedTVList for query. It should clone TVList if necessary. + Map alignedTvListQueryMap = + prepareTvListMapForQuery( + context, alignedMemChunk, modsToMemtable == null, globalTimeFilter, columnIndexList); + if (modsToMemtable != null) { timeColumnDeletion = ModificationUtils.constructDeletionList( @@ -696,7 +827,7 @@ public ReadOnlyMemChunk getReadOnlyMemChunkFromMemTable( } // prepare TVList for query. It should clone TVList if necessary. Map tvListQueryMap = - prepareTvListMapForQuery(context, memChunk, modsToMemtable == null, globalTimeFilter); + prepareTvListMapForQuery(context, memChunk, modsToMemtable == null, globalTimeFilter, null); List 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 9eb30ec59119e..f138324150e48 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 @@ -26,6 +26,7 @@ import org.apache.iotdb.db.queryengine.execution.fragment.FragmentInstanceContext; import org.apache.iotdb.db.queryengine.execution.fragment.QueryContext; 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; @@ -43,6 +44,7 @@ import java.io.UncheckedIOException; import java.util.Iterator; import java.util.List; +import java.util.Set; import java.util.concurrent.BlockingQueue; public abstract class AbstractWritableMemChunk implements IWritableMemChunk { @@ -109,19 +111,39 @@ protected void maybeReleaseTvList(TVList tvList) { } } + /** + * 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 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 0ba04ad53e683..55e1899ceeb22 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 @@ -128,12 +128,12 @@ public void sortTvLists() { // 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); @@ -387,12 +387,12 @@ public IPointReader getPointReader() { 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 39629bbbcaae2..bfefb9817bb37 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 @@ -139,11 +139,11 @@ public void sortTvLists() { 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); @@ -298,11 +298,11 @@ public IPointReader getPointReader() { 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 1e29d245756b0..7ca1d474fc619 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 @@ -23,6 +23,7 @@ import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory; import org.apache.iotdb.db.i18n.DataNodeMiscMessages; import org.apache.iotdb.db.i18n.StorageEngineMessages; +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; @@ -58,8 +59,10 @@ 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; @@ -87,6 +90,56 @@ public abstract class AlignedTVList extends TVList { private long materializedBitmapMemoryCost; private long arrayMemCostWithoutPrimitiveArraysAndIndex; + /** + * 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[] valueColumnsToMove; + private final List[] 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[] valueColumnsToMove, + List[] 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> values; @@ -110,12 +163,14 @@ public abstract class AlignedTVList extends TVList { dataTypes = types; memoryBinaryChunkSize = new long[dataTypes.size()]; materializedValueArrayCounts = new int[dataTypes.size()]; - refreshArrayMemCostWithoutPrimitiveArrays(); values = new ArrayList<>(types.size()); for (int i = 0; i < types.size(); i++) { values.add(new ArrayList<>(getDefaultArrayNum())); } + // arrayMemCostWithoutPrimitiveArrays depends on per-column value arrays, so values must be + // initialized before computing it + refreshArrayMemCostWithoutPrimitiveArrays(); } public static AlignedTVList newAlignedList(List dataTypes) { @@ -177,7 +232,7 @@ public TVList getTvListByColumnIndex( alignedTvList.allValueColDeletedMap = ignoreAllNullRows ? getAllValueColDeletedMap() : null; alignedTvList.timeColDeletedMap = this.timeColDeletedMap; alignedTvList.timeDeletedCnt = this.timeDeletedCnt; - alignedTvList.materializedBitmapMemoryCost = calculateBitmapRamCost(bitMaps); + alignedTvList.materializedBitmapMemoryCost = calculateBitmapRamCost(bitMaps, null); for (int i = 0; i < columnIndexList.size(); i++) { int columnIndex = columnIndexList.get(i); if (columnIndex != -1 && values.get(i) != null) { @@ -187,7 +242,6 @@ public TVList getTvListByColumnIndex( (long) materializedArrayCount * valueListArrayMemCost(dataTypeList.get(i)); } } - return alignedTvList; } @@ -212,38 +266,168 @@ public synchronized AlignedTVList clone() { AlignedTVList cloneList = AlignedTVList.newAlignedList(new ArrayList<>(dataTypes)); cloneAs(cloneList); cloneList.timeDeletedCnt = this.timeDeletedCnt; - System.arraycopy( - memoryBinaryChunkSize, 0, cloneList.memoryBinaryChunkSize, 0, dataTypes.size()); - for (int i = 0; i < values.size(); i++) { - // Clone value + cloneColumnDataTo(cloneList, null); + cloneList.timeColDeletedMap = timeColDeletedMap == null ? null : timeColDeletedMap.clone(); + cloneList.materializedValueArrayCounts = + Arrays.copyOf(materializedValueArrayCounts, materializedValueArrayCounts.length); + cloneList.materializedValueArrayMemCost = materializedValueArrayMemCost; + 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 columnsToClone) { + Set 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 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[] valueColumnsToMove = (List[]) new List[columnCount]; + List[] bitmapColumnsToMove = (List[]) new List[columnCount]; + for (int i = 0; i < columnCount; i++) { + if (retainedColumns.contains(i)) { + continue; + } + List 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 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 == 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)); } - if (cloneList.bitMaps.get(i) == null) { - List cloneColumnBitMaps = new ArrayList<>(columnBitMaps.size()); - for (BitMap bitMap : columnBitMaps) { - cloneColumnBitMaps.add(bitMap == null ? null : bitMap.clone()); - } - cloneList.bitMaps.set(i, cloneColumnBitMaps); + bitmapColumnsToMove[i] = bitMaps.get(i); + } + } + + return new PartialClonePlan( + this, + cloneList, + valueColumnsToMove, + bitmapColumnsToMove, + calculateArrayMemCostWithoutIndex(retainedColumns), + cloneList.calculateArrayMemCostWithoutIndex(null), + calculateBitmapRamCost(bitMaps, retainedColumns), + calculateBitmapRamCost(bitMaps, null)); + } + + private synchronized void commitPartialClone(PartialClonePlan plan) { + // The clone keeps the deep-copied retained columns, so copy their accounting too. + for (int i = 0; i < dataTypes.size(); i++) { + int materializedCount = materializedValueArrayCounts[i]; + plan.cloneList.materializedValueArrayCounts[i] = materializedCount; + if (materializedCount > 0) { + plan.cloneList.materializedValueArrayMemCost += + (long) materializedCount * valueListArrayMemCost(dataTypes.get(i)); + } + } + for (int i = 0; i < plan.valueColumnsToMove.length; i++) { + List columnValues = plan.valueColumnsToMove[i]; + if (columnValues == null) { + continue; + } + + // Move value arrays and bitmaps from the source to the clone. The clone was created with + // empty column containers, so the moved references are set into place. + plan.cloneList.values.set(i, columnValues); + values.set(i, null); + List columnBitMaps = plan.bitmapColumnsToMove[i]; + if (columnBitMaps != null) { + plan.cloneList.bitMaps.set(i, columnBitMaps); + bitMaps.set(i, null); + } + memoryBinaryChunkSize[i] = 0; + + // The source no longer owns the moved column's materialized arrays. + int materializedCount = materializedValueArrayCounts[i]; + materializedValueArrayCounts[i] = 0; + if (materializedCount > 0) { + materializedValueArrayMemCost -= + (long) materializedCount * valueListArrayMemCost(dataTypes.get(i)); + } + } + + materializedBitmapMemoryCost = plan.sourceBitmapMemoryCost; + plan.cloneList.materializedBitmapMemoryCost = plan.cloneBitmapMemoryCost; + // Refresh the cached per-block cost after the retained columns changed on both lists. + refreshArrayMemCostWithoutPrimitiveArrays(); + plan.cloneList.refreshArrayMemCostWithoutPrimitiveArrays(); + } + + /** + * 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 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 columnValues = values.get(i); + if (columnValues == null) { + continue; + } + + // Release memory for non-query columns + for (Object dataArray : columnValues) { + if (dataArray != null) { + 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); + } + + // Remove the released column from the materialized-array accounting + int materializedCount = materializedValueArrayCounts[i]; + materializedValueArrayCounts[i] = 0; + if (materializedCount > 0) { + materializedValueArrayMemCost -= + (long) materializedCount * valueListArrayMemCost(dataTypes.get(i)); + } } - cloneList.timeColDeletedMap = timeColDeletedMap == null ? null : timeColDeletedMap.clone(); - cloneList.materializedValueArrayCounts = - Arrays.copyOf(materializedValueArrayCounts, materializedValueArrayCounts.length); - cloneList.materializedValueArrayMemCost = materializedValueArrayMemCost; - cloneList.materializedBitmapMemoryCost = materializedBitmapMemoryCost; - return cloneList; + + materializedBitmapMemoryCost = calculateBitmapRamCost(bitMaps, columnsToKeep); + // Refresh the cached per-block cost after releasing columns + refreshArrayMemCostWithoutPrimitiveArrays(); } @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity warning @@ -641,6 +825,25 @@ public List getTsDataTypes() { 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 getAccessedColumnsForQuery() { + Set 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 @@ -745,6 +948,13 @@ public synchronized Pair delete( * or delete wrong rows. */ public synchronized void deleteColumn(int columnIndex) { + List 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> localBitMaps = new ArrayList<>(dataTypes.size()); for (int j = 0; j < dataTypes.size(); j++) { @@ -752,10 +962,11 @@ public synchronized void deleteColumn(int columnIndex) { } bitMaps = localBitMaps; } + if (bitMaps.get(columnIndex) == null) { - List columnBitMaps = new ArrayList<>(values.get(columnIndex).size()); - for (int i = 0; i < values.get(columnIndex).size(); i++) { - columnBitMaps.add(BitMap.createBitMapDynamically(ARRAY_SIZE)); + List columnBitMaps = new ArrayList<>(columnValues.size()); + for (int i = 0; i < columnValues.size(); i++) { + columnBitMaps.add(new BitMap(ARRAY_SIZE)); } bitMaps.set(columnIndex, columnBitMaps); materializedBitmapMemoryCost += @@ -815,6 +1026,69 @@ protected Object cloneValue(TSDataType type, Object value) { } } + /* + * 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 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 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 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 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++) { @@ -852,7 +1126,12 @@ protected void expandValues() { indices.add((int[]) getPrimitiveArraysByType(TSDataType.INT32)); } for (int i = 0; i < dataTypes.size(); i++) { - values.get(i).add(null); + List 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(null); if (bitMaps != null && bitMaps.get(i) != null) { bitMaps.get(i).add(null); materializedBitmapMemoryCost += bitmapReferenceRamCost(); @@ -1132,6 +1411,14 @@ private Object getOrCreateValueArray(int columnIndex, int arrayIndex) { } private BitMap getBitMap(int columnIndex, int arrayIndex) { + List 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> localBitMaps = new ArrayList<>(dataTypes.size()); @@ -1143,8 +1430,8 @@ private BitMap getBitMap(int columnIndex, int arrayIndex) { // if the bitmap in columnIndex is null, init the bitmap of this column from the beginning if (bitMaps.get(columnIndex) == null) { - List columnBitMaps = new ArrayList<>(values.get(columnIndex).size()); - for (int i = 0; i < values.get(columnIndex).size(); i++) { + List columnBitMaps = new ArrayList<>(columnValues.size()); + for (int i = 0; i < columnValues.size(); i++) { columnBitMaps.add(null); } bitMaps.set(columnIndex, columnBitMaps); @@ -1186,22 +1473,59 @@ public synchronized RamInfo calculateRamSize() { new ArrayList<>(dataTypes)); } + public synchronized RamInfo calculateRamSize(Set columnsToClone) { + return new RamInfo( + timestamps.size(), + alignedTvListArrayMemCost(columnsToClone), + getRamSize(columnsToClone), + rowCount, + new ArrayList<>(dataTypes)); + } + public synchronized long getRamSize() { return (long) timestamps.size() * alignedTvListArrayMemCostWithoutPrimitiveArrays() + materializedValueArrayMemCost + materializedBitmapMemoryCost; } - private static long calculateBitmapRamCost(List> bitMaps) { + public synchronized long getRamSize(Set columnsToClone) { + long size = + (long) timestamps.size() * alignedTvListArrayMemCostWithoutPrimitiveArrays(columnsToClone); + for (int i = 0; i < dataTypes.size(); i++) { + if (columnsToClone != null && !columnsToClone.contains(i)) { + continue; + } + TSDataType dataType = dataTypes.get(i); + if (dataType != null) { + size += (long) materializedValueArrayCounts[i] * valueListArrayMemCost(dataType); + } + } + return size + calculateBitmapRamCost(bitMaps, columnsToClone); + } + + private long calculateArrayMemCostWithoutIndex(Set retainedColumns) { + long arrayMemCost = alignedTvListArrayMemCost(retainedColumns); + if (indices != null) { + arrayMemCost -= (long) PrimitiveArrayManager.ARRAY_SIZE * Integer.BYTES; + } + return arrayMemCost; + } + + private static long calculateBitmapRamCost( + List> bitMaps, Set columnsToClone) { if (bitMaps == null) { return 0; } long size = 0; - for (List columnBitMaps : bitMaps) { + for (int i = 0, length = bitMaps.size(); i < length; i++) { + if (columnsToClone != null && !columnsToClone.contains(i)) { + continue; + } + List 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 += bitMap.ramBytesUsed(); @@ -1249,27 +1573,29 @@ public static long alignedTvListArrayMemCost( * * @return AlignedTvListArrayMemSize */ - public long alignedTvListArrayMemCost() { + public long alignedTvListArrayMemCost(Set columnsToClone) { long size = 0; - // value & bitmap array mem size + 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; } @@ -1278,13 +1604,30 @@ public long alignedTvListArrayMemCostWithoutPrimitiveArrays() { + (indices != null ? (long) PrimitiveArrayManager.ARRAY_SIZE * Integer.BYTES : 0); } + private long alignedTvListArrayMemCostWithoutPrimitiveArrays(Set retainedColumns) { + long size = alignedTvListArrayMemCost(retainedColumns); + if (indices != null) { + size -= (long) PrimitiveArrayManager.ARRAY_SIZE * Integer.BYTES; + } + for (int i = 0; i < dataTypes.size(); i++) { + TSDataType dataType = dataTypes.get(i); + if (dataType != null + && values.get(i) != null + && (retainedColumns == null || retainedColumns.contains(i))) { + size -= valueListArrayMemCost(dataType); + } + } + return size; + } + private void refreshArrayMemCostWithoutPrimitiveArrays() { long size = alignedTvListArrayMemCost(); if (indices != null) { size -= (long) PrimitiveArrayManager.ARRAY_SIZE * Integer.BYTES; } - for (TSDataType dataType : dataTypes) { - if (dataType != null) { + for (int i = 0; i < dataTypes.size(); i++) { + TSDataType dataType = dataTypes.get(i); + if (dataType != null && values.get(i) != null) { size -= valueListArrayMemCost(dataType); } } @@ -1304,6 +1647,10 @@ public static long alignedTvListArrayMemCostWithoutPrimitiveArrays( return size; } + public long alignedTvListArrayMemCost() { + return alignedTvListArrayMemCost((Set) 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 9dbd11cf6268d..fbd593097a112 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 @@ -824,6 +824,16 @@ public Set getQueryContextSet() { 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 getAccessedColumnsForQuery() { + return null; + } + public List 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 e8c7994fb2013..51dff0c8d52a2 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 @@ -50,6 +50,7 @@ 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.BeforeClass; import org.junit.Test; @@ -316,4 +317,24 @@ private IMemTable createMemTable(String deviceId, String measurementId) } return memTable; } + + private IMemTable createMemTable(String deviceId, List 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 6d0cabb044313..5f0a1ad8f9ea9 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.calc.exception.MemoryNotEnoughException; import org.apache.iotdb.db.queryengine.plan.planner.memory.NotThreadSafeMemoryReservationManager; import org.junit.After; @@ -161,4 +162,49 @@ public void testMemoryReservationManagerNormalPriorityReserveAndRelease() { 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 0000000000000..cd1ece64f083c --- /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> result = null; + try { + temporaryList.lockQueryList(); + try { + result = + executor.submit( + () -> + resourceByPathUtils.prepareTvListMapForQuery( + currentQuery, memChunk, false, null, null)); + Assert.assertTrue(temporaryListInitialized.await(3, TimeUnit.SECONDS)); + Future> blockedResult = result; + Assert.assertThrows( + TimeoutException.class, () -> blockedResult.get(200, TimeUnit.MILLISECONDS)); + } finally { + temporaryList.unlockQueryList(); + } + + Map 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/utils/datastructure/AlignedTVListTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java index 381f733bf4dd7..4dca43170fc5d 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 @@ -36,7 +36,10 @@ import java.nio.ByteBuffer; 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; @@ -497,4 +500,162 @@ public void testCalculateChunkSize() { Assert.assertEquals(tvList.memoryBinaryChunkSize[0], 0); Assert.assertEquals(tvList.memoryBinaryChunkSize[1], 0); } + + @Test + public void testMovesUnclonedColumns() { + List 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 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 testPartialRamSizeScalesWithRetainedColumns() { + int columnCount = 256; + List 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 retainedColumns = Collections.singleton(0); + long retainedRamSize = tvList.calculateRamSize(retainedColumns).getRamSize(); + long fullRamSize = tvList.calculateRamSize().getRamSize(); + + // Only the retained column's materialized arrays are charged, so keeping 1 of 256 columns + // must cost far less than the full list. + Assert.assertTrue(retainedRamSize < fullRamSize / 64); + + 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 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 firstColumnValues = tvList.getValues().get(0); + List secondColumnValues = tvList.getValues().get(1); + List thirdColumnValues = tvList.getValues().get(2); + List firstColumnBitMaps = tvList.getBitMaps().get(0); + Object invalidThirdColumnArray = new int[ARRAY_SIZE]; + thirdColumnValues.set(0, invalidThirdColumnArray); + + Set 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 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 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)); + } + } } From 82801d816f5a58a3be9358b88fed9ff65acfa767 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Fri, 7 Aug 2026 12:33:37 +0800 Subject: [PATCH 2/6] fix(test): implement releaseMemoryImmediately in MemoryReservationManager test fakes The cherry-pick added the releaseMemoryImmediately(long) method to the calc-commons MemoryReservationManager interface, but four test classes that implement the interface directly were not updated, breaking test compilation. Add the missing method to each fake, mirroring the semantics of the existing releaseMemoryCumulatively implementation in each class. --- .../grouped/rate/GroupedRateAccumulatorMemoryTest.java | 5 +++++ .../aggregation/rate/RateAccumulatorFactoryTest.java | 3 +++ .../rate/RateFunctionIntermediateStateCodecTest.java | 5 +++++ .../execution/fragment/QueryModificationLoaderTest.java | 5 +++++ 4 files changed, 18 insertions(+) diff --git a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/grouped/rate/GroupedRateAccumulatorMemoryTest.java b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/grouped/rate/GroupedRateAccumulatorMemoryTest.java index 7aa15ebfdea3f..54c0dc4478ab8 100644 --- a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/grouped/rate/GroupedRateAccumulatorMemoryTest.java +++ b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/grouped/rate/GroupedRateAccumulatorMemoryTest.java @@ -74,6 +74,11 @@ public void releaseMemoryCumulatively(long size) { cumulativeRelease += size; } + @Override + public void releaseMemoryImmediately(long size) { + cumulativeRelease += size; + } + @Override public void releaseAllReservedMemory() {} diff --git a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateAccumulatorFactoryTest.java b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateAccumulatorFactoryTest.java index 364e4439c2681..e2556a3cc03df 100644 --- a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateAccumulatorFactoryTest.java +++ b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateAccumulatorFactoryTest.java @@ -161,6 +161,9 @@ public void reserveMemoryImmediately(long size) {} @Override public void releaseMemoryCumulatively(long size) {} + @Override + public void releaseMemoryImmediately(long size) {} + @Override public void releaseAllReservedMemory() {} diff --git a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateFunctionIntermediateStateCodecTest.java b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateFunctionIntermediateStateCodecTest.java index ea535e05c5c25..aed2a17d96c37 100644 --- a/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateFunctionIntermediateStateCodecTest.java +++ b/iotdb-core/calc-commons/src/test/java/org/apache/iotdb/calc/execution/operator/source/relational/aggregation/rate/RateFunctionIntermediateStateCodecTest.java @@ -128,6 +128,11 @@ public void releaseMemoryCumulatively(long size) { outstandingReservation -= size; } + @Override + public void releaseMemoryImmediately(long size) { + outstandingReservation -= size; + } + @Override public void releaseAllReservedMemory() { outstandingReservation = 0; diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/QueryModificationLoaderTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/QueryModificationLoaderTest.java index 65cbc0f5b9f2b..14673ffda6465 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/QueryModificationLoaderTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/fragment/QueryModificationLoaderTest.java @@ -323,6 +323,11 @@ public void releaseMemoryCumulatively(long size) { reservedBytes -= size; } + @Override + public void releaseMemoryImmediately(long size) { + reservedBytes -= size; + } + @Override public void releaseAllReservedMemory() { reservedBytes = 0; From 02b8dde3cfe5089dd77c01d9c747ef20bbd6de8c Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Fri, 7 Aug 2026 14:03:52 +0800 Subject: [PATCH 3/6] spotless --- .../db/schemaengine/schemaregion/utils/ResourceByPathUtils.java | 2 +- .../plan/planner/LocalExecutionPlannerOperatorsMemoryTest.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) 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 8c7793896c736..f95fcfda473b7 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 @@ -74,8 +74,8 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; -import java.util.stream.Collectors; import java.util.Set; +import java.util.stream.Collectors; import static org.apache.iotdb.commons.path.AlignedPath.VECTOR_PLACEHOLDER; 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 5f0a1ad8f9ea9..f736250183370 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 @@ -19,8 +19,8 @@ package org.apache.iotdb.db.queryengine.plan.planner; -import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.calc.exception.MemoryNotEnoughException; +import org.apache.iotdb.db.queryengine.common.QueryId; import org.apache.iotdb.db.queryengine.plan.planner.memory.NotThreadSafeMemoryReservationManager; import org.junit.After; From dff9170ec9dc369af27daf34530aaac9335107b9 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Fri, 7 Aug 2026 14:12:57 +0800 Subject: [PATCH 4/6] fix(test): adapt cherry-picked tests to master path/device-id APIs - ResourceByPathUtilsTest: getResourceInstance(IFullPath) requires a NonAlignedFullPath; MeasurementPath does not implement IFullPath, so the test no longer compiled. Build the NonAlignedFullPath from IDeviceID.Factory.DEFAULT_FACTORY and a MeasurementSchema. - FragmentInstanceExecutionTest: use IDeviceID.Factory.DEFAULT_FACTORY.create directly instead of DeviceIDFactory.getInstance().getDeviceID(new PartialPath(...)), matching the query-side NonAlignedFullPath construction in the same file, and drop the now-unused imports. --- .../execution/fragment/FragmentInstanceExecutionTest.java | 6 ++---- .../schemaregion/utils/ResourceByPathUtilsTest.java | 8 ++++++-- 2 files changed, 8 insertions(+), 6 deletions(-) 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 51dff0c8d52a2..d43f7b6c46bdf 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 @@ -25,7 +25,6 @@ import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.path.NonAlignedFullPath; -import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.queryengine.common.FragmentInstanceId; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; @@ -35,7 +34,6 @@ import org.apache.iotdb.db.queryengine.execution.exchange.sink.ISink; import org.apache.iotdb.db.queryengine.execution.schedule.IDriverScheduler; import org.apache.iotdb.db.storageengine.dataregion.DataRegion; -import org.apache.iotdb.db.storageengine.dataregion.memtable.DeviceIDFactory; import org.apache.iotdb.db.storageengine.dataregion.memtable.IMemTable; import org.apache.iotdb.db.storageengine.dataregion.memtable.IWritableMemChunk; import org.apache.iotdb.db.storageengine.dataregion.memtable.IWritableMemChunkGroup; @@ -309,7 +307,7 @@ private IMemTable createMemTable(String deviceId, String measurementId) int rows = 100; for (int i = 0; i < 100; i++) { memTable.write( - DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)), + IDeviceID.Factory.DEFAULT_FACTORY.create(deviceId), Collections.singletonList( new MeasurementSchema(measurementId, TSDataType.INT32, TSEncoding.PLAIN)), rows - i - 1, @@ -330,7 +328,7 @@ private IMemTable createMemTable(String deviceId, List schem values[j] = (long) i * 100 + j; } memTable.writeAlignedRow( - DeviceIDFactory.getInstance().getDeviceID(new PartialPath(deviceId)), + IDeviceID.Factory.DEFAULT_FACTORY.create(deviceId), schemaList, i, values); 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 index cd1ece64f083c..334e6836e7968 100644 --- 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 @@ -18,12 +18,14 @@ */ package org.apache.iotdb.db.schemaengine.schemaregion.utils; -import org.apache.iotdb.commons.path.MeasurementPath; +import org.apache.iotdb.commons.path.NonAlignedFullPath; 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.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.write.schema.MeasurementSchema; import org.junit.Assert; import org.junit.Test; @@ -70,7 +72,9 @@ public void testFlushingQueryLocksTemporaryTVListBeforeRegistration() throws Exc ResourceByPathUtils resourceByPathUtils = ResourceByPathUtils.getResourceInstance( - new MeasurementPath("root.test.d.s", TSDataType.INT64)); + new NonAlignedFullPath( + IDeviceID.Factory.DEFAULT_FACTORY.create("root.test.d"), + new MeasurementSchema("s", TSDataType.INT64))); QueryContext currentQuery = new QueryContext(2, false); ExecutorService executor = Executors.newSingleThreadExecutor(); Future> result = null; From 7cfb464b4223bfa1e830c42dbc08f246838d9c46 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Fri, 7 Aug 2026 15:54:32 +0800 Subject: [PATCH 5/6] spotless --- .../execution/fragment/FragmentInstanceExecutionTest.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) 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 d43f7b6c46bdf..2332cde925419 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 @@ -328,10 +328,7 @@ private IMemTable createMemTable(String deviceId, List schem values[j] = (long) i * 100 + j; } memTable.writeAlignedRow( - IDeviceID.Factory.DEFAULT_FACTORY.create(deviceId), - schemaList, - i, - values); + IDeviceID.Factory.DEFAULT_FACTORY.create(deviceId), schemaList, i, values); } return memTable; } From aa4ab1106684e57ff7053ac398697139f6d01946 Mon Sep 17 00:00:00 2001 From: shuwenwei Date: Fri, 7 Aug 2026 18:19:12 +0800 Subject: [PATCH 6/6] fix bug --- .../db/utils/datastructure/AlignedTVList.java | 26 +++++-------------- 1 file changed, 6 insertions(+), 20 deletions(-) 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 cdbdf60959346..f85b996ca4e85 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 @@ -100,8 +100,6 @@ public static final class PartialClonePlan { private final AlignedTVList cloneList; private final List[] valueColumnsToMove; private final List[] bitmapColumnsToMove; - private final long sourceArrayMemCostWithoutIndex; - private final long cloneArrayMemCostWithoutIndex; private final long sourceBitmapMemoryCost; private final long cloneBitmapMemoryCost; @@ -112,16 +110,12 @@ private PartialClonePlan( AlignedTVList cloneList, List[] valueColumnsToMove, List[] 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; } @@ -266,9 +260,7 @@ public synchronized AlignedTVList cloneForFlushSort() { public synchronized AlignedTVList clone() { AlignedTVList cloneList = AlignedTVList.newAlignedList(new ArrayList<>(dataTypes)); cloneAs(cloneList); - cloneList.timeDeletedCnt = this.timeDeletedCnt; cloneColumnDataTo(cloneList, null); - cloneList.timeColDeletedMap = timeColDeletedMap == null ? null : timeColDeletedMap.clone(); cloneList.materializedValueArrayCounts = Arrays.copyOf(materializedValueArrayCounts, materializedValueArrayCounts.length); cloneList.materializedValueArrayMemCost = materializedValueArrayMemCost; @@ -331,8 +323,6 @@ private PartialClonePlan prepareMovePlan(AlignedTVList cloneList, Set r cloneList, valueColumnsToMove, bitmapColumnsToMove, - calculateArrayMemCostWithoutIndex(retainedColumns), - cloneList.calculateArrayMemCostWithoutIndex(null), calculateBitmapRamCost(bitMaps, retainedColumns), calculateBitmapRamCost(bitMaps, null)); } @@ -990,12 +980,16 @@ protected Object cloneValue(TSDataType type, Object value) { * 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 + * This method only performs the allocation phase: copy row-level time deletion state, 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 columnsToClone) { + cloneList.timeDeletedCnt = timeDeletedCnt; + cloneList.timeColDeletedMap = timeColDeletedMap == null ? null : timeColDeletedMap.clone(); + boolean cloneAllColumns = columnsToClone == null; System.arraycopy( memoryBinaryChunkSize, 0, cloneList.memoryBinaryChunkSize, 0, dataTypes.size()); @@ -1446,14 +1440,6 @@ public synchronized long getRamSize(Set columnsToClone) { return size + calculateBitmapRamCost(bitMaps, columnsToClone); } - private long calculateArrayMemCostWithoutIndex(Set retainedColumns) { - long arrayMemCost = alignedTvListArrayMemCost(retainedColumns); - if (indices != null) { - arrayMemCost -= (long) PrimitiveArrayManager.ARRAY_SIZE * Integer.BYTES; - } - return arrayMemCost; - } - private static long calculateBitmapRamCost( List> bitMaps, Set columnsToClone) { if (bitMaps == null) {