This is an automated email from the ASF dual-hosted git repository. JackieTien97 pushed a commit to branch rc/2.0.11 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit b17bf5e39b39c036b1c380995f25c94cba76aa3b Author: Caideyipi <[email protected]> AuthorDate: Mon Jul 20 16:30:20 2026 +0800 [Performance] Reduce aligned MemTable bitmap memory usage (#18249) --- .../dataregion/memtable/TsFileProcessor.java | 24 +++++---- .../db/utils/datastructure/AlignedTVList.java | 61 ++++++++++++++++++---- .../dataregion/memtable/TsFileProcessorTest.java | 28 +++++----- .../db/utils/datastructure/AlignedTVListTest.java | 51 ++++++++++++++++++ 4 files changed, 129 insertions(+), 35 deletions(-) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java index 6dc5ecb3e2b..ab2066c1c6a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java @@ -824,8 +824,10 @@ public class TsFileProcessor { } // this insertion will result in a new array if ((alignedMemChunk.alignedListSize() % PrimitiveArrayManager.ARRAY_SIZE) == 0) { - dataTypesInTVList.addAll(alignedMemChunk.getWorkingTVList().getTsDataTypes()); memTableIncrement += alignedMemChunk.getWorkingTVList().alignedTvListArrayMemCost(); + for (TSDataType dataType : dataTypesInTVList) { + memTableIncrement += AlignedTVList.valueListArrayMemCost(dataType); + } } } @@ -913,13 +915,13 @@ public class TsFileProcessor { int addingPointNum = addingPointNumInfo.right; // Here currentChunkPointNum + addingPointNum >= 1 if (((currentChunkPointNum + addingPointNum) % PrimitiveArrayManager.ARRAY_SIZE) == 0) { - if (alignedMemChunk != null) { - dataTypesInTVList.addAll(alignedMemChunk.getWorkingTVList().getTsDataTypes()); - } dataTypesInTVList.addAll(addingPointNumInfo.left.values()); memTableIncrement += alignedMemChunk != null ? alignedMemChunk.getWorkingTVList().alignedTvListArrayMemCost() + + dataTypesInTVList.stream() + .mapToLong(AlignedTVList::valueListArrayMemCost) + .sum() : AlignedTVList.alignedTvListArrayMemCost( dataTypesInTVList.toArray(new TSDataType[0]), null); } @@ -1077,6 +1079,9 @@ public class TsFileProcessor { List<TSDataType> dataTypesInTVList = new ArrayList<>(); int currentPointNum = alignedMemChunk.alignedListSize(); int newPointNum = currentPointNum + incomingPointNum; + int currentArrayCnt = + currentPointNum / PrimitiveArrayManager.ARRAY_SIZE + + (currentPointNum % PrimitiveArrayManager.ARRAY_SIZE > 0 ? 1 : 0); for (int i = 0; dataTypes != null && i < dataTypes.length; i++) { TSDataType dataType = dataTypes[i]; if (!isWritableFieldMeasurement(measurementIds, dataTypes, columns, columnCategories, i)) { @@ -1085,17 +1090,12 @@ public class TsFileProcessor { if (!alignedMemChunk.containsMeasurement(measurementIds[i])) { // add a new column in the TVList, the new column should be as long as existing ones - memIncrements[0] += - (currentPointNum / PrimitiveArrayManager.ARRAY_SIZE + 1) - * AlignedTVList.valueListArrayMemCost(dataType); + memIncrements[0] += currentArrayCnt * AlignedTVList.valueListArrayMemCost(dataType); dataTypesInTVList.add(dataType); } } // calculate how many new arrays will be added after this insertion - int currentArrayCnt = - currentPointNum / PrimitiveArrayManager.ARRAY_SIZE - + (currentPointNum % PrimitiveArrayManager.ARRAY_SIZE > 0 ? 1 : 0); int newArrayCnt = newPointNum / PrimitiveArrayManager.ARRAY_SIZE + (newPointNum % PrimitiveArrayManager.ARRAY_SIZE > 0 ? 1 : 0); @@ -1103,9 +1103,11 @@ public class TsFileProcessor { if (acquireArray != 0) { // memory of extending the TVList - dataTypesInTVList.addAll(alignedMemChunk.getWorkingTVList().getTsDataTypes()); memIncrements[0] += acquireArray * alignedMemChunk.getWorkingTVList().alignedTvListArrayMemCost(); + for (TSDataType dataType : dataTypesInTVList) { + memIncrements[0] += acquireArray * AlignedTVList.valueListArrayMemCost(dataType); + } } } 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 01af929d5c8..0e44a3dbde5 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 @@ -46,6 +46,7 @@ import org.apache.tsfile.read.filter.basic.Filter; import org.apache.tsfile.utils.Binary; import org.apache.tsfile.utils.BitMap; import org.apache.tsfile.utils.Pair; +import org.apache.tsfile.utils.RamUsageEstimator; import org.apache.tsfile.utils.ReadWriteForEncodingUtils; import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.utils.TsPrimitiveType; @@ -72,6 +73,11 @@ import static org.apache.tsfile.utils.RamUsageEstimator.NUM_BYTES_OBJECT_REF; public abstract class AlignedTVList extends TVList { + private static final long BITMAP_RAM_COST_PER_BLOCK = + RamUsageEstimator.shallowSizeOfInstance(BitMap.class) + + RamUsageEstimator.sizeOfByteArray(BitMap.getSizeOfBytes(ARRAY_SIZE)) + + NUM_BYTES_OBJECT_REF; + // Data types of this aligned tvList protected List<TSDataType> dataTypes; @@ -925,7 +931,9 @@ public abstract class AlignedTVList extends TVList { /* 1. Build result-level bitmap (1 = failure row) */ byte[] resultBitMap = - (results != null) ? buildResultBitMapBytes(results, idx, elementIdx, len) : null; + results != null && containsFailedStatus(results, idx, len) + ? buildResultBitMapBytes(results, idx, elementIdx, len) + : null; for (int j = 0; j < values.length; j++) { /* Fast-path: column is entirely null */ @@ -935,7 +943,7 @@ public abstract class AlignedTVList extends TVList { } /* 2.mask the column bitmap */ - if (bitMaps != null && bitMaps[j] != null) { + if (bitMaps != null && bitMaps[j] != null && containsMarkedBit(bitMaps[j], idx, len)) { getBitMap(j, arrayIndex).merge(bitMaps[j], idx, elementIdx, len); } @@ -946,12 +954,46 @@ public abstract class AlignedTVList extends TVList { } } + private static boolean containsFailedStatus(TSStatus[] results, int start, int length) { + for (int i = start; i < start + length; i++) { + if (results[i] != null && results[i].code != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { + return true; + } + } + return false; + } + + private static boolean containsMarkedBit(BitMap bitMap, int start, int length) { + if (length <= 0) { + return false; + } + + byte[] bytes = bitMap.getByteArray(); + int end = start + length - 1; + int firstByteIndex = start >>> 3; + int lastByteIndex = end >>> 3; + if (firstByteIndex == lastByteIndex) { + int mask = (0xFF << (start & 7)) & (0xFF >>> (7 - (end & 7))); + return (bytes[firstByteIndex] & mask) != 0; + } + + if ((bytes[firstByteIndex] & (0xFF << (start & 7))) != 0) { + return true; + } + for (int i = firstByteIndex + 1; i < lastByteIndex; i++) { + if (bytes[i] != 0) { + return true; + } + } + return (bytes[lastByteIndex] & (0xFF >>> (7 - (end & 7)))) != 0; + } + public static byte[] buildResultBitMapBytes( TSStatus[] results, int idx, int elementIdx, int length) { int start = elementIdx & 7; int totalBits = start + length; int size = (totalBits + 7) >> 3; - BitMap bitmap = new BitMap(size, new byte[size]); + BitMap bitmap = new BitMap(totalBits, new byte[size]); if (results == null) { return bitmap.getByteArray(); @@ -1028,14 +1070,14 @@ public abstract class AlignedTVList extends TVList { if (bitMaps.get(columnIndex) == null) { List<BitMap> columnBitMaps = new ArrayList<>(values.get(columnIndex).size()); for (int i = 0; i < values.get(columnIndex).size(); i++) { - columnBitMaps.add(new BitMap(ARRAY_SIZE, new byte[ARRAY_SIZE])); + columnBitMaps.add(null); } bitMaps.set(columnIndex, columnBitMaps); } // if the bitmap in arrayIndex is null, init the bitmap if (bitMaps.get(columnIndex).get(arrayIndex) == null) { - bitMaps.get(columnIndex).set(arrayIndex, new BitMap(ARRAY_SIZE, new byte[ARRAY_SIZE])); + bitMaps.get(columnIndex).set(arrayIndex, new BitMap(ARRAY_SIZE)); } return bitMaps.get(columnIndex).get(arrayIndex); @@ -1089,6 +1131,7 @@ public abstract class AlignedTVList extends TVList { if (type != null && (columnCategories == null || columnCategories[i] == TsTableColumnCategory.FIELD)) { size += (long) ARRAY_SIZE * (long) type.getDataTypeSize(); + size += BITMAP_RAM_COST_PER_BLOCK; measurementColumnNum++; } } @@ -1117,9 +1160,7 @@ public abstract class AlignedTVList extends TVList { TSDataType type = dataTypes.get(column); if (type != null) { size += (long) PrimitiveArrayManager.ARRAY_SIZE * (long) type.getDataTypeSize(); - if (bitMaps != null && bitMaps.get(column) != null) { - size += (long) PrimitiveArrayManager.ARRAY_SIZE / 8 + 1; - } + size += BITMAP_RAM_COST_PER_BLOCK; } } // size is 0 when all types are null @@ -1147,8 +1188,8 @@ public abstract class AlignedTVList extends TVList { long size = 0; // value array mem size size += (long) PrimitiveArrayManager.ARRAY_SIZE * (long) type.getDataTypeSize(); - // bitmap array mem size - size += (long) PrimitiveArrayManager.ARRAY_SIZE / 8 + 1; + // bitmap object, byte array, and reference in the bitmap list + size += BITMAP_RAM_COST_PER_BLOCK; // array headers mem size size += NUM_BYTES_ARRAY_HEADER; // Object references size in ArrayList diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java index 2b14116ff2d..040988997b2 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessorTest.java @@ -536,21 +536,21 @@ public class TsFileProcessorTest { true, new long[5]); IMemTable memTable = processor.getWorkMemTable(); - Assert.assertEquals(1596552, memTable.getTVListsRamCost()); + Assert.assertEquals(1776552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNode(100, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(1596552, memTable.getTVListsRamCost()); + Assert.assertEquals(1776552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNode(200, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(1596552, memTable.getTVListsRamCost()); + Assert.assertEquals(1776552, memTable.getTVListsRamCost()); Assert.assertEquals(90000, memTable.getTotalPointsNum()); Assert.assertEquals(720360, memTable.memSize()); // Test records @@ -559,7 +559,7 @@ public class TsFileProcessorTest { record.addTuple(DataPoint.getDataPoint(dataType, measurementId, String.valueOf(i))); processor.insert(buildInsertRowNodeByTSRecord(record), new long[5]); } - Assert.assertEquals(1598168, memTable.getTVListsRamCost()); + Assert.assertEquals(1778168, memTable.getTVListsRamCost()); Assert.assertEquals(90100, memTable.getTotalPointsNum()); Assert.assertEquals(721560, memTable.memSize()); } @@ -614,56 +614,56 @@ public class TsFileProcessorTest { true, new long[5]); IMemTable memTable = processor.getWorkMemTable(); - Assert.assertEquals(1596552, memTable.getTVListsRamCost()); + Assert.assertEquals(1776552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNodeFors3000ToS6000(0, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(3219552, memTable.getTVListsRamCost()); + Assert.assertEquals(3552552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNode(100, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(3219552, memTable.getTVListsRamCost()); + Assert.assertEquals(3552552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNodeFors3000ToS6000(100, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(3219552, memTable.getTVListsRamCost()); + Assert.assertEquals(3552552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNode(200, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(3219552, memTable.getTVListsRamCost()); + Assert.assertEquals(3552552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNodeFors3000ToS6000(200, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(3219552, memTable.getTVListsRamCost()); + Assert.assertEquals(3552552, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNode(300, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(6466104, memTable.getTVListsRamCost()); + Assert.assertEquals(7105104, memTable.getTVListsRamCost()); processor.insertTablet( genInsertTableNodeFors3000ToS6000(300, true), Collections.singletonList(new int[] {0, 10}), new TSStatus[10], true, new long[5]); - Assert.assertEquals(6466104, memTable.getTVListsRamCost()); + Assert.assertEquals(7105104, memTable.getTVListsRamCost()); Assert.assertEquals(240000, memTable.getTotalPointsNum()); Assert.assertEquals(1920960, memTable.memSize()); @@ -673,14 +673,14 @@ public class TsFileProcessorTest { record.addTuple(DataPoint.getDataPoint(dataType, measurementId, String.valueOf(i))); processor.insert(buildInsertRowNodeByTSRecord(record), new long[5]); } - Assert.assertEquals(6467720, memTable.getTVListsRamCost()); + Assert.assertEquals(7106720, memTable.getTVListsRamCost()); // Test records for (int i = 1; i <= 100; i++) { TSRecord record = new TSRecord(deviceId, i); record.addTuple(DataPoint.getDataPoint(dataType, "s1", String.valueOf(i))); processor.insert(buildInsertRowNodeByTSRecord(record), new long[5]); } - Assert.assertEquals(6469336, memTable.getTVListsRamCost()); + Assert.assertEquals(7108336, memTable.getTVListsRamCost()); Assert.assertEquals(240200, memTable.getTotalPointsNum()); Assert.assertEquals(1923360, memTable.memSize()); } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/utils/datastructure/AlignedTVListTest.java index 10d7803a83f..5ab976c1310 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 @@ -18,6 +18,9 @@ */ package org.apache.iotdb.db.utils.datastructure; +import org.apache.iotdb.common.rpc.thrift.TSStatus; +import org.apache.iotdb.rpc.TSStatusCode; + import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.external.commons.lang3.ArrayUtils; @@ -27,8 +30,11 @@ import org.junit.Assert; import org.junit.Test; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; +import static org.apache.iotdb.db.storageengine.rescon.memory.PrimitiveArrayManager.ARRAY_SIZE; + public class AlignedTVListTest { @Test @@ -142,6 +148,51 @@ public class AlignedTVListTest { } } + @Test + public void testBitmapIsAllocatedLazilyWithCompactBackingArray() { + AlignedTVList tvList = + AlignedTVList.newAlignedList(Arrays.asList(TSDataType.INT64, TSDataType.INT64)); + Object[] values = new Object[] {1L, 1L}; + for (int i = 0; i < ARRAY_SIZE * 2 + 1; i++) { + tvList.putAlignedValue(i, values); + } + + Assert.assertNull(tvList.getBitMaps()); + tvList.putAlignedValue(ARRAY_SIZE * 2 + 1L, new Object[] {null, 1L}); + + List<BitMap> firstColumnBitMaps = tvList.getBitMaps().get(0); + Assert.assertEquals(3, firstColumnBitMaps.size()); + Assert.assertNull(firstColumnBitMaps.get(0)); + Assert.assertNull(firstColumnBitMaps.get(1)); + Assert.assertNotNull(firstColumnBitMaps.get(2)); + Assert.assertEquals( + BitMap.getSizeOfBytes(ARRAY_SIZE), firstColumnBitMaps.get(2).getByteArray().length); + Assert.assertTrue(tvList.isNullValue(ARRAY_SIZE * 2 + 1, 0)); + Assert.assertFalse(tvList.isNullValue(ARRAY_SIZE * 2, 0)); + } + + @Test + public void testEmptyInputBitmapsDoNotMaterializeMemTableBitmaps() { + AlignedTVList tvList = AlignedTVList.newAlignedList(List.of(TSDataType.INT64)); + long[] times = new long[ARRAY_SIZE]; + long[][] values = new long[1][ARRAY_SIZE]; + BitMap[] bitMaps = new BitMap[] {new BitMap(ARRAY_SIZE)}; + TSStatus[] results = new TSStatus[ARRAY_SIZE]; + for (int i = 0; i < ARRAY_SIZE; i++) { + times[i] = i; + values[0][i] = i; + } + + tvList.putAlignedValues(times, values, bitMaps, 0, ARRAY_SIZE, results); + + Assert.assertNull(tvList.getBitMaps()); + + Arrays.fill(results, new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); + tvList.putAlignedValues(times, values, bitMaps, 0, ARRAY_SIZE, results); + + Assert.assertNull(tvList.getBitMaps()); + } + @Test public void testClone() { List<TSDataType> dataTypes = new ArrayList<>();
