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/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/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..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 @@ -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; @@ -73,6 +74,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.stream.Collectors; 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 035b9a89bb92e..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 @@ -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,50 @@ 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 sourceBitmapMemoryCost; + private final long cloneBitmapMemoryCost; + + private boolean committed; + + private PartialClonePlan( + AlignedTVList sourceList, + AlignedTVList cloneList, + List[] valueColumnsToMove, + List[] bitmapColumnsToMove, + long sourceBitmapMemoryCost, + long cloneBitmapMemoryCost) { + this.sourceList = sourceList; + this.cloneList = cloneList; + this.valueColumnsToMove = valueColumnsToMove; + this.bitmapColumnsToMove = bitmapColumnsToMove; + 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. // A null primitive array means all existing rows in that block are null for the column. // Index relation: columnIndex(dataTypeIndex) -> arrayIndex -> elementIndex @@ -111,12 +158,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) { @@ -178,7 +227,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) { @@ -188,7 +237,6 @@ public TVList getTvListByColumnIndex( (long) materializedArrayCount * valueListArrayMemCost(dataTypeList.get(i)); } } - return alignedTvList; } @@ -212,39 +260,165 @@ public synchronized AlignedTVList cloneForFlushSort() { 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.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, + 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 @@ -618,6 +792,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 @@ -778,6 +971,73 @@ 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: 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()); + 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++) { @@ -815,7 +1075,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(); @@ -1079,6 +1344,14 @@ private Object getOrCreateValueArray(int columnIndex, int arrayIndex, int elemen } 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()); @@ -1090,8 +1363,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); @@ -1137,22 +1410,51 @@ 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 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(); @@ -1200,27 +1502,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; } @@ -1229,13 +1533,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); } } @@ -1255,6 +1576,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..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 @@ -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; @@ -50,6 +48,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; @@ -308,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, @@ -316,4 +315,21 @@ 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( + IDeviceID.Factory.DEFAULT_FACTORY.create(deviceId), schemaList, i, values); + } + return memTable; + } } 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; 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..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,6 +19,7 @@ package org.apache.iotdb.db.queryengine.plan.planner; +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; @@ -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..334e6836e7968 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/schemaengine/schemaregion/utils/ResourceByPathUtilsTest.java @@ -0,0 +1,113 @@ +/* + * 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.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; + +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 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; + 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 fbbd35cff4316..d9db4db6cc55b 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; @@ -586,4 +589,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)); + } + } }