diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java index 5fab6979e2912..3bebbc9115305 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java @@ -94,6 +94,8 @@ public class PipeInsertNodeTabletInsertionEvent extends PipeInsertionEvent private final AtomicReference allocatedMemoryBlock; private volatile List tablets; + // Calculated together with tablets so downstream batching does not rescan Tablet internals. + private volatile long tabletsMemoryUsageInBytes; private List eventParsers; @@ -481,22 +483,31 @@ public boolean isAligned(final int i) { // TODO: for table model insertion, we need to get the database name public synchronized List convertToTablets() { if (Objects.isNull(tablets)) { - tablets = - initEventParsers().stream() - .map(TabletInsertionEventParser::convertToTablet) - .collect(Collectors.toList()); + final List parsers = initEventParsers(); + final List convertedTablets = new ArrayList<>(parsers.size()); + long tabletMemoryUsageInBytes = 0; + for (final TabletInsertionEventParser parser : parsers) { + final Tablet tablet = parser.convertToTablet(); + convertedTablets.add(tablet); + // Tablet.ramBytesUsed() is required for the memory block to account for the actual + // retained tablet size. Calculate it while converting to avoid a second stream traversal. + tabletMemoryUsageInBytes += PipeMemoryWeightUtil.calculateTabletSizeInBytes(tablet); + } + tablets = convertedTablets; + tabletsMemoryUsageInBytes = tabletMemoryUsageInBytes; allocatedMemoryBlock.compareAndSet( null, PipeDataNodeResourceManager.memory() - .forceAllocateForTabletWithRetry( - tablets.stream() - .map(PipeMemoryWeightUtil::calculateTabletSizeInBytes) - .reduce(Long::sum) - .orElse(0L))); + .forceAllocateForTabletWithRetry(tabletMemoryUsageInBytes)); } return tablets; } + public long getTabletsMemoryUsageInBytes() { + convertToTablets(); + return tabletsMemoryUsageInBytes; + } + /////////////////////////// event parser /////////////////////////// private List initEventParsers() { @@ -505,35 +516,27 @@ private List initEventParsers() { return eventParsers; } - eventParsers = new ArrayList<>(); final InsertNode node = getInsertNode(); if (Objects.isNull(node)) { throw new PipeException(DataNodePipeMessages.INSERTNODE_HAS_BEEN_RELEASED); } + eventParsers = new ArrayList<>(getEventParserCount(node)); + final UserEntity userEntity = + shouldParse4Privilege + ? new UserEntity(Long.parseLong(userId), userName, cliHostname) + : null; switch (node.getType()) { case INSERT_ROW: case INSERT_TABLET: eventParsers.add( new TabletInsertionEventTreePatternParser( - pipeTaskMeta, - this, - node, - treePattern, - shouldParse4Privilege - ? new UserEntity(Long.parseLong(userId), userName, cliHostname) - : null)); + pipeTaskMeta, this, node, treePattern, userEntity)); break; case INSERT_ROWS: for (final InsertRowNode insertRowNode : ((InsertRowsNode) node).getInsertRowNodeList()) { eventParsers.add( new TabletInsertionEventTreePatternParser( - pipeTaskMeta, - this, - insertRowNode, - treePattern, - shouldParse4Privilege - ? new UserEntity(Long.parseLong(userId), userName, cliHostname) - : null)); + pipeTaskMeta, this, insertRowNode, treePattern, userEntity)); } break; case RELATIONAL_INSERT_ROW: @@ -565,6 +568,16 @@ private List initEventParsers() { } } + private static int getEventParserCount(final InsertNode node) { + if (node instanceof InsertRowsNode) { + return ((InsertRowsNode) node).getInsertRowNodeList().size(); + } + if (node instanceof RelationalInsertRowsNode) { + return ((RelationalInsertRowsNode) node).getInsertRowNodeList().size(); + } + return 1; + } + public long count() { long count = 0; for (final Tablet covertedTablet : convertToTablets()) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java index 5b50eb166bebb..384f475a3592b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java @@ -58,12 +58,12 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.HashMap; import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Objects; import java.util.function.Predicate; -import java.util.stream.Collectors; public class TsFileInsertionEventTableParserTabletIterator implements Iterator { @@ -137,9 +137,12 @@ public TsFileInsertionEventTableParserTabletIterator( this.metadataQuerier = new MetadataQuerierByFileImpl(reader); fileMetadata = this.metadataQuerier.getWholeFileMetadata(); final List> tableSchemaList = - fileMetadata.getTableSchemaMap().entrySet().stream() - .filter(predicate) - .collect(Collectors.toList()); + new ArrayList<>(fileMetadata.getTableSchemaMap().size()); + for (final Map.Entry entry : fileMetadata.getTableSchemaMap().entrySet()) { + if (predicate.test(entry)) { + tableSchemaList.add(entry); + } + } this.allocatedMemoryBlockForTablet = allocatedMemoryBlockForTablet; this.allocatedMemoryBlockForBatchData = allocatedMemoryBlockForBatchData; @@ -250,10 +253,10 @@ public boolean hasNext() { deviceMetaIterator = metadataQuerier.deviceIterator(tableRoot, null); final int columnSchemaSize = tableSchema.getColumnSchemas().size(); - dataTypeList = new ArrayList<>(); - columnTypes = new ArrayList<>(); - measurementList = new ArrayList<>(); - fieldSchemaList = new ArrayList<>(); + dataTypeList = new ArrayList<>(columnSchemaSize); + columnTypes = new ArrayList<>(columnSchemaSize); + measurementList = new ArrayList<>(columnSchemaSize); + fieldSchemaList = new ArrayList<>(columnSchemaSize); for (int i = 0; i < columnSchemaSize; i++) { final IMeasurementSchema schema = tableSchema.getColumnSchemas().get(i); @@ -364,28 +367,27 @@ private void initChunkReader(final AbstractAlignedChunkMetadata alignedChunkMeta timeChunk.getData().rewind(); long size = timeChunkSize; - final List valueChunkList = new ArrayList<>(); + final int fieldSchemaSize = fieldSchemaList.size(); + final List valueChunkList = new ArrayList<>(fieldSchemaSize); final Map valueChunkMetadataMap = - alignedChunkMetadata.getValueChunkMetadataList().stream() - .filter(Objects::nonNull) - .filter( - metadata -> - !isFieldDeletedByMods( - metadata.getMeasurementUid(), - alignedChunkMetadata.getStartTime(), - alignedChunkMetadata.getEndTime())) - .collect( - Collectors.toMap( - IChunkMetadata::getMeasurementUid, - metadata -> metadata, - (left, right) -> left)); + new HashMap<>((int) (fieldSchemaSize / 0.75f) + 1); + for (final IChunkMetadata metadata : alignedChunkMetadata.getValueChunkMetadataList()) { + if (metadata != null + && !isFieldDeletedByMods( + metadata.getMeasurementUid(), + alignedChunkMetadata.getStartTime(), + alignedChunkMetadata.getEndTime())) { + // Keep the first metadata entry to preserve the former merge-function behavior. + valueChunkMetadataMap.putIfAbsent(metadata.getMeasurementUid(), metadata); + } + } // To ensure that the Tablet has the same alignedChunk column as the current one, // you need to create a new Tablet to fill in the data. isSameDeviceID = false; // Need to ensure that columnTypes recreates an array - final List categories = new ArrayList<>(deviceIdSize); + final List categories = new ArrayList<>(deviceIdSize + fieldSchemaSize); for (int i = 0; i < deviceIdSize; i++) { categories.add(ColumnCategory.TAG); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/InsertNodeMemoryEstimator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/InsertNodeMemoryEstimator.java index 54655704c1c53..9e4e1c2fa8f81 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/InsertNodeMemoryEstimator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/InsertNodeMemoryEstimator.java @@ -733,24 +733,16 @@ private static long sizeOfObjectList(final List list) { if (list == null) { return 0L; } - long size = RamUsageEstimator.shallowSizeOf(list); - if (list instanceof ArrayList) { - size += - RamUsageEstimator.alignObjectSize( - NUM_BYTES_ARRAY_HEADER + NUM_BYTES_OBJECT_REF * list.size()); - } - return size; + return SIZE_OF_ARRAYLIST + + RamUsageEstimator.alignObjectSize( + NUM_BYTES_ARRAY_HEADER + NUM_BYTES_OBJECT_REF * list.size()); } private static long sizeOfIntegerList(final List integers) { if (integers == null) { return 0L; } - long size = sizeOfObjectList(integers); - for (Integer ignored : integers) { - size += SIZE_OF_INT; - } - return size; + return sizeOfObjectList(integers) + (long) SIZE_OF_INT * integers.size(); } private static long sizeOfResults(final Map results) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventPlainBatch.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventPlainBatch.java index bdf6ee1874a4d..3eec34cc94bd8 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventPlainBatch.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventPlainBatch.java @@ -131,7 +131,8 @@ public PipeTransferTabletBatchReqV2 toTPipeTransferReq() throws IOException { } } for (final Pair tabletPair : batchTablets) { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateTabletSerializedSize(tabletPair.getRight())); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { tabletPair.getRight().serialize(outputStream); ReadWriteIOUtils.write(true, outputStream); @@ -176,7 +177,12 @@ private long buildTabletInsertionBuffer(final TabletInsertionEvent event) throws insertNodeDataBases.add(databaseName); } else { final List tablets = pipeInsertNodeTabletInsertionEvent.convertToTablets(); - estimateSize = calculateTabletsSizeInBytes(tablets); + // convertToTablets() has already measured every tablet for the event memory block. Reuse + // that exact measurement instead of calling Tablet.ramBytesUsed() (which walks the schema + // map) once more while building this batch. + estimateSize = + pipeInsertNodeTabletInsertionEvent.getTabletsMemoryUsageInBytes() + + (long) Integer.BYTES * tablets.size(); increaseTotalBufferSizeAndUpdateMemoryBlock(estimateSize); for (final Tablet tablet : tablets) { constructTabletBatchWithoutMemoryReservation( @@ -192,9 +198,11 @@ private long buildTabletInsertionBuffer(final TabletInsertionEvent event) throws pipeRawTabletInsertionEvent.convertToTablet(), pipeRawTabletInsertionEvent.getTableModelDatabaseName()); } else { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + final Tablet tablet = pipeRawTabletInsertionEvent.convertToTablet(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateTabletSerializedSize(tablet)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { - pipeRawTabletInsertionEvent.convertToTablet().serialize(outputStream); + tablet.serialize(outputStream); ReadWriteIOUtils.write(pipeRawTabletInsertionEvent.isAligned(), outputStream); buffer = ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); } @@ -226,14 +234,14 @@ private void constructTabletBatchWithoutMemoryReservation( currentBatch.getRight().add(tablet); } - private long calculateTabletsSizeInBytes(final List tablets) { - return tablets.stream().mapToLong(PipeTabletEventPlainBatch::calculateTabletSizeInBytes).sum(); - } - private static long calculateTabletSizeInBytes(final Tablet tablet) { return PipeMemoryWeightUtil.calculateTabletSizeInBytes(tablet) + 4; } + private static int calculateTabletSerializedSize(final Tablet tablet) { + return tablet.serializedSize() + Byte.BYTES; + } + static boolean mayAppendTablet(final Tablet target, final Tablet source) { // Tablet.append already checks schemas and column categories. Avoid repeating those potentially // expensive comparisons here because wide-table pipe transfer can have many columns. diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReq.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReq.java index 352ff0bfc63a2..4de77179a7cf0 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReq.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReq.java @@ -46,6 +46,11 @@ public class PipeTransferTabletBatchReq extends TPipeTransferReq { + private static final int BATCH_REQUEST_COUNT_SERIALIZED_SIZE = + Integer.BYTES // legacy binary request count + + Integer.BYTES // insert node request count + + Integer.BYTES; // raw tablet request count + private final transient List binaryReqs = new ArrayList<>(); private final transient List insertNodeReqs = new ArrayList<>(); private final transient List tabletReqs = new ArrayList<>(); @@ -129,19 +134,28 @@ public static PipeTransferTabletBatchReq toTPipeTransferReq( batchReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); batchReq.type = PipeRequestType.TRANSFER_TABLET_BATCH.getType(); - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateSerializedSize(insertNodeBuffers, tabletBuffers)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { // Binary buffer, for rolling upgrade ReadWriteIOUtils.write(0, outputStream); + // Insert-node and raw-tablet serializations are self-delimiting, so their lengths are not + // written separately. ReadWriteIOUtils.write(insertNodeBuffers.size(), outputStream); for (final ByteBuffer insertNodeBuffer : insertNodeBuffers) { - outputStream.write(insertNodeBuffer.array(), 0, insertNodeBuffer.limit()); + outputStream.write( + insertNodeBuffer.array(), + insertNodeBuffer.arrayOffset() + insertNodeBuffer.position(), + insertNodeBuffer.remaining()); } ReadWriteIOUtils.write(tabletBuffers.size(), outputStream); for (final ByteBuffer tabletBuffer : tabletBuffers) { - outputStream.write(tabletBuffer.array(), 0, tabletBuffer.limit()); + outputStream.write( + tabletBuffer.array(), + tabletBuffer.arrayOffset() + tabletBuffer.position(), + tabletBuffer.remaining()); } batchReq.body = @@ -151,6 +165,13 @@ public static PipeTransferTabletBatchReq toTPipeTransferReq( return batchReq; } + static int calculateSerializedSize( + final List insertNodeBuffers, final List tabletBuffers) { + return BATCH_REQUEST_COUNT_SERIALIZED_SIZE + + insertNodeBuffers.stream().mapToInt(ByteBuffer::remaining).sum() + + tabletBuffers.stream().mapToInt(ByteBuffer::remaining).sum(); + } + public static PipeTransferTabletBatchReq fromTPipeTransferReq( final TPipeTransferReq transferReq) { final PipeTransferTabletBatchReq batchReq = new PipeTransferTabletBatchReq(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java index 6c4607518b4e7..1ef33fa32c83b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java @@ -46,6 +46,12 @@ import java.util.Objects; public class PipeTransferTabletBatchReqV2 extends TPipeTransferReq { + + private static final int BATCH_REQUEST_COUNT_SERIALIZED_SIZE = + Integer.BYTES // legacy binary request count + + Integer.BYTES // insert node request count + + Integer.BYTES; // raw tablet request count + private final transient List insertNodeReqs = new ArrayList<>(); private final transient List tabletReqs = new ArrayList<>(); @@ -55,7 +61,8 @@ private PipeTransferTabletBatchReqV2() { } public List constructStatements() { - final List statements = new ArrayList<>(); + final List statements = + new ArrayList<>(insertNodeReqs.size() + tabletReqs.size()); final Map> tableModelDatabaseInsertRowStatementMap = new LinkedHashMap<>(); @@ -186,22 +193,33 @@ public static PipeTransferTabletBatchReqV2 toTPipeTransferReq( batchReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); batchReq.type = PipeRequestType.TRANSFER_TABLET_BATCH_V2.getType(); - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS( + calculateSerializedSize( + insertNodeBuffers, tabletBuffers, insertNodeDataBases, tabletDataBases)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { // Binary buffer, for rolling upgrade ReadWriteIOUtils.write(0, outputStream); + // Insert-node and raw-tablet serializations are self-delimiting, so their lengths are not + // written separately. ReadWriteIOUtils.write(insertNodeBuffers.size(), outputStream); for (int i = 0; i < insertNodeBuffers.size(); i++) { final ByteBuffer insertNodeBuffer = insertNodeBuffers.get(i); - outputStream.write(insertNodeBuffer.array(), 0, insertNodeBuffer.limit()); + outputStream.write( + insertNodeBuffer.array(), + insertNodeBuffer.arrayOffset() + insertNodeBuffer.position(), + insertNodeBuffer.remaining()); ReadWriteIOUtils.write(insertNodeDataBases.get(i), outputStream); } ReadWriteIOUtils.write(tabletBuffers.size(), outputStream); for (int i = 0; i < tabletBuffers.size(); i++) { final ByteBuffer tabletBuffer = tabletBuffers.get(i); - outputStream.write(tabletBuffer.array(), 0, tabletBuffer.limit()); + outputStream.write( + tabletBuffer.array(), + tabletBuffer.arrayOffset() + tabletBuffer.position(), + tabletBuffer.remaining()); ReadWriteIOUtils.write(tabletDataBases.get(i), outputStream); } @@ -212,6 +230,23 @@ public static PipeTransferTabletBatchReqV2 toTPipeTransferReq( return batchReq; } + static int calculateSerializedSize( + final List insertNodeBuffers, + final List tabletBuffers, + final List insertNodeDataBases, + final List tabletDataBases) { + int size = BATCH_REQUEST_COUNT_SERIALIZED_SIZE; + for (int i = 0; i < insertNodeBuffers.size(); i++) { + size += insertNodeBuffers.get(i).remaining(); + size += ReadWriteIOUtils.sizeToWrite(insertNodeDataBases.get(i)); + } + for (int i = 0; i < tabletBuffers.size(); i++) { + size += tabletBuffers.get(i).remaining(); + size += ReadWriteIOUtils.sizeToWrite(tabletDataBases.get(i)); + } + return size; + } + public static PipeTransferTabletBatchReqV2 fromTPipeTransferReq( final org.apache.iotdb.service.rpc.thrift.TPipeTransferReq transferReq) { final PipeTransferTabletBatchReqV2 batchReq = new PipeTransferTabletBatchReqV2(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReq.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReq.java index b31816c1fcd6a..cf4a746cb0eaf 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReq.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReq.java @@ -32,7 +32,6 @@ import org.apache.iotdb.db.storageengine.dataregion.wal.buffer.WALEntry; import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq; -import org.apache.tsfile.utils.BytesUtils; import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; @@ -102,14 +101,23 @@ public static PipeTransferTabletBinaryReq fromTPipeTransferReq( /////////////////////////////// Air Gap /////////////////////////////// public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer) throws IOException { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateSerializedSize(byteBuffer)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_BINARY.getType(), outputStream); - return BytesUtils.concatByteArray(byteArrayOutputStream.toByteArray(), byteBuffer.array()); + outputStream.write( + byteBuffer.array(), + byteBuffer.arrayOffset() + byteBuffer.position(), + byteBuffer.remaining()); + return byteArrayOutputStream.toByteArray(); } } + static int calculateSerializedSize(final ByteBuffer byteBuffer) { + return Byte.BYTES + Short.BYTES + byteBuffer.remaining(); + } + /////////////////////////////// Object /////////////////////////////// @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java index 2788033be2d04..196af0b16d2d1 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java @@ -119,10 +119,14 @@ public static PipeTransferTabletBinaryReqV2 toTPipeTransferReq( req.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); req.type = PipeRequestType.TRANSFER_TABLET_BINARY_V2.getType(); - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateSerializedSize(byteBuffer, dataBaseName)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { - ReadWriteIOUtils.write(byteBuffer.limit(), outputStream); - outputStream.write(byteBuffer.array(), 0, byteBuffer.limit()); + ReadWriteIOUtils.write(byteBuffer.remaining(), outputStream); + outputStream.write( + byteBuffer.array(), + byteBuffer.arrayOffset() + byteBuffer.position(), + byteBuffer.remaining()); ReadWriteIOUtils.write(dataBaseName, outputStream); req.body = ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); } @@ -151,17 +155,29 @@ public static PipeTransferTabletBinaryReqV2 fromTPipeTransferReq( public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer, final String dataBaseName) throws IOException { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateAirGapSerializedSize(byteBuffer, dataBaseName)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_BINARY_V2.getType(), outputStream); - ReadWriteIOUtils.write(byteBuffer.limit(), outputStream); - outputStream.write(byteBuffer.array(), 0, byteBuffer.limit()); + ReadWriteIOUtils.write(byteBuffer.remaining(), outputStream); + outputStream.write( + byteBuffer.array(), + byteBuffer.arrayOffset() + byteBuffer.position(), + byteBuffer.remaining()); ReadWriteIOUtils.write(dataBaseName, outputStream); return byteArrayOutputStream.toByteArray(); } } + static int calculateSerializedSize(final ByteBuffer byteBuffer, final String dataBaseName) { + return Integer.BYTES + byteBuffer.remaining() + ReadWriteIOUtils.sizeToWrite(dataBaseName); + } + + static int calculateAirGapSerializedSize(final ByteBuffer byteBuffer, final String dataBaseName) { + return Byte.BYTES + Short.BYTES + calculateSerializedSize(byteBuffer, dataBaseName); + } + /////////////////////////////// Object /////////////////////////////// @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java index bc42630d79b4d..42f353c250667 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java @@ -31,7 +31,6 @@ import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement; import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq; -import org.apache.tsfile.utils.BytesUtils; import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; @@ -106,15 +105,28 @@ public static PipeTransferTabletInsertNodeReq fromTPipeTransferReq( /////////////////////////////// Air Gap /////////////////////////////// public static byte[] toTPipeTransferBytes(final InsertNode insertNode) throws IOException { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateAirGapSerializedSize(insertNode)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_INSERT_NODE.getType(), outputStream); - return BytesUtils.concatByteArray( - byteArrayOutputStream.toByteArray(), insertNode.serializeToByteBuffer().array()); + insertNode.serialize(outputStream); + return byteArrayOutputStream.toByteArray(); } } + static int calculateSerializedSize(final InsertNode insertNode) { + return insertNode.serializeToByteBufferSize(); + } + + static int calculateAirGapSerializedSize(final InsertNode insertNode) { + return calculateAirGapSerializedSize(calculateSerializedSize(insertNode)); + } + + protected static int calculateAirGapSerializedSize(final int bodySize) { + return Byte.BYTES + Short.BYTES + bodySize; + } + /////////////////////////////// Object /////////////////////////////// @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java index b9d5eb7de85bf..4fc18398dda75 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReqV2.java @@ -120,7 +120,8 @@ public static PipeTransferTabletInsertNodeReqV2 toTPipeTransferReq( req.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); req.type = PipeRequestType.TRANSFER_TABLET_INSERT_NODE_V2.getType(); - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateSerializedSize(insertNode, dataBaseName)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { insertNode.serialize(outputStream); ReadWriteIOUtils.write(req.dataBaseName, outputStream); @@ -150,7 +151,8 @@ public static PipeTransferTabletInsertNodeReqV2 fromTPipeTransferReq( public static byte[] toTPipeTransferBytes(final InsertNode insertNode, final String dataBaseName) throws IOException { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateAirGapSerializedSize(insertNode, dataBaseName)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write( @@ -161,6 +163,16 @@ public static byte[] toTPipeTransferBytes(final InsertNode insertNode, final Str } } + static int calculateSerializedSize(final InsertNode insertNode, final String dataBaseName) { + return PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode) + + ReadWriteIOUtils.sizeToWrite(dataBaseName); + } + + static int calculateAirGapSerializedSize(final InsertNode insertNode, final String dataBaseName) { + return PipeTransferTabletInsertNodeReq.calculateAirGapSerializedSize( + calculateSerializedSize(insertNode, dataBaseName)); + } + /////////////////////////////// Object /////////////////////////////// @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReq.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReq.java index 01c80758152d7..00907ed8302f5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReq.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReq.java @@ -135,7 +135,7 @@ public static PipeTransferTabletRawReq toTPipeTransferReq( tabletReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); tabletReq.type = PipeRequestType.TRANSFER_TABLET_RAW.getType(); - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(calculateSerializedSize(tablet)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { tablet.serialize(outputStream); ReadWriteIOUtils.write(isAligned, outputStream); @@ -272,7 +272,8 @@ public byte[] toTPipeTransferBytes() throws IOException { throw new IOException(DataNodePipeMessages.CANNOT_SERIALIZE_BOTH_TABLET_AND_STATEMENT_ARE); } - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateAirGapSerializedSize(tabletToSerialize)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_RAW.getType(), outputStream); @@ -298,6 +299,14 @@ public static byte[] toTPipeTransferBytes(final Tablet tablet, final boolean isA return req.toTPipeTransferBytes(); } + static int calculateSerializedSize(final Tablet tablet) { + return tablet.serializedSize() + Byte.BYTES; + } + + static int calculateAirGapSerializedSize(final Tablet tablet) { + return Byte.BYTES + Short.BYTES + calculateSerializedSize(tablet); + } + /////////////////////////////// Object /////////////////////////////// @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReqV2.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReqV2.java index d395bf6cf5f26..0c827261fc22d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReqV2.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReqV2.java @@ -144,7 +144,8 @@ public static PipeTransferTabletRawReqV2 toTPipeTransferReq( tabletReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); tabletReq.type = PipeRequestType.TRANSFER_TABLET_RAW_V2.getType(); - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateSerializedSize(tablet, dataBaseName)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { tablet.serialize(outputStream); ReadWriteIOUtils.write(isAligned, outputStream); @@ -173,7 +174,8 @@ public static PipeTransferTabletRawReqV2 fromTPipeTransferReq( public static byte[] toTPipeTransferBytes( final Tablet tablet, final boolean isAligned, final String dataBaseName) throws IOException { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateAirGapSerializedSize(tablet, dataBaseName)); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_RAW_V2.getType(), outputStream); @@ -184,6 +186,14 @@ public static byte[] toTPipeTransferBytes( } } + static int calculateSerializedSize(final Tablet tablet, final String dataBaseName) { + return tablet.serializedSize() + Byte.BYTES + ReadWriteIOUtils.sizeToWrite(dataBaseName); + } + + static int calculateAirGapSerializedSize(final Tablet tablet, final String dataBaseName) { + return Byte.BYTES + Short.BYTES + calculateSerializedSize(tablet, dataBaseName); + } + /////////////////////////////// Object /////////////////////////////// @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/builder/IoTConsensusV2TransferBatchReqBuilder.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/builder/IoTConsensusV2TransferBatchReqBuilder.java index fc387084a0012..8c9e0299f31ff 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/builder/IoTConsensusV2TransferBatchReqBuilder.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/builder/IoTConsensusV2TransferBatchReqBuilder.java @@ -40,7 +40,6 @@ import org.slf4j.LoggerFactory; import java.io.IOException; -import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; import java.util.List; @@ -209,7 +208,6 @@ public List deepCopyEvents() { } protected int buildTabletInsertionBuffer(TabletInsertionEvent event) throws WALPipeException { - final ByteBuffer buffer; final TCommitId commitId; // event instanceof PipeInsertNodeTabletInsertionEvent) @@ -221,17 +219,15 @@ protected int buildTabletInsertionBuffer(TabletInsertionEvent event) throws WALP pipeInsertNodeTabletInsertionEvent.getCommitterKey().getRestartTimes(), pipeInsertNodeTabletInsertionEvent.getRebootTimes()); - // Read the bytebuffer from the wal file and transfer it directly without serializing or - // deserializing if possible final InsertNode insertNode = pipeInsertNodeTabletInsertionEvent.getInsertNode(); // IoTConsensusV2 will transfer binary data to TIoTConsensusV2TransferReq final ProgressIndex progressIndex = pipeInsertNodeTabletInsertionEvent.getProgressIndex(); - buffer = insertNode.serializeToByteBuffer(); - batchReqs.add( + final IoTConsensusV2TabletInsertNodeReq request = IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq( - insertNode, commitId, consensusGroupId, progressIndex, thisDataNodeId)); + insertNode, commitId, consensusGroupId, progressIndex, thisDataNodeId); + batchReqs.add(request); - return buffer.limit(); + return request.body.remaining(); } @Override diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/request/IoTConsensusV2TabletInsertNodeReq.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/request/IoTConsensusV2TabletInsertNodeReq.java index 5f076b68ec387..af1d0a9fc5f34 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/request/IoTConsensusV2TabletInsertNodeReq.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/request/IoTConsensusV2TabletInsertNodeReq.java @@ -95,6 +95,7 @@ public static IoTConsensusV2TabletInsertNodeReq toTIoTConsensusV2TransferReq( req.dataNodeId = thisDataNodeId; req.version = IoTConsensusV2RequestVersion.VERSION_1.getVersion(); req.type = IoTConsensusV2RequestType.TRANSFER_TABLET_INSERT_NODE.getType(); + // InsertNode preallocates this buffer with its manually calculated Pipe serialization size. req.body = insertNode.serializeToByteBuffer(); try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(); @@ -109,6 +110,11 @@ public static IoTConsensusV2TabletInsertNodeReq toTIoTConsensusV2TransferReq( return req; } + /** Returns the exact serialized size of an InsertNode request body. */ + public static int calculateSerializedSize(final InsertNode insertNode) { + return insertNode.serializeToByteBufferSize(); + } + public static IoTConsensusV2TabletInsertNodeReq fromTIoTConsensusV2TransferReq( TIoTConsensusV2TransferReq transferReq) { final IoTConsensusV2TabletInsertNodeReq insertNodeReq = new IoTConsensusV2TabletInsertNodeReq(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java index a1be7a44fa7e2..d39281f286432 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java @@ -251,6 +251,10 @@ public void clearUselessField() { typeProvider = null; } + public void clearUselessFieldsAfterRouting() { + planNodeTree.clearUselessFieldsAfterRouting(); + } + public void clearTypeProvider() { typeProvider = null; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/pipe/PipeEnrichedInsertNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/pipe/PipeEnrichedInsertNode.java index f6c323a1b5759..0f527ab45424c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/pipe/PipeEnrichedInsertNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/pipe/PipeEnrichedInsertNode.java @@ -39,6 +39,7 @@ import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.write.schema.MeasurementSchema; import java.io.DataOutputStream; @@ -163,6 +164,11 @@ public void setDataRegionReplicaSet(final TRegionReplicaSet dataRegionReplicaSet insertNode.setDataRegionReplicaSet(dataRegionReplicaSet); } + @Override + public void clearUselessFieldsAfterRouting() { + insertNode.clearUselessFieldsAfterRouting(); + } + @Override public PartialPath getTargetPath() { return insertNode.getTargetPath(); @@ -290,6 +296,16 @@ protected void serializeAttributes(final DataOutputStream stream) throws IOExcep insertNode.serialize(stream); } + @Override + protected int serializedAttributesSize() { + return PlanNodeType.BYTES + insertNode.serializeToByteBufferSize(); + } + + @Override + protected int serializedPlanNodeIdSize() { + return ReadWriteIOUtils.sizeToWrite(super.getPlanNodeId().getId()); + } + public static PipeEnrichedInsertNode deserialize(final ByteBuffer buffer) { return new PipeEnrichedInsertNode((InsertNode) PlanNodeType.deserialize(buffer)); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertMultiTabletsNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertMultiTabletsNode.java index 4d1b987b89272..8d484933e8f3d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertMultiTabletsNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertMultiTabletsNode.java @@ -280,6 +280,15 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { } } + @Override + protected int serializedAttributesSize() { + int size = PlanNodeType.BYTES + Integer.BYTES; + for (final InsertTabletNode insertTabletNode : insertTabletNodeList) { + size += insertTabletNode.serializedSubAttributesSize(); + } + return size + parentInsertTabletNodeIndexList.size() * Integer.BYTES; + } + @Override public void markAsGeneratedByPipe() { isGeneratedByPipe = true; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java index 9c8e369188357..289866c6882f2 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java @@ -22,6 +22,7 @@ import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; import org.apache.iotdb.commons.consensus.index.ProgressIndex; import org.apache.iotdb.commons.exception.IllegalPathException; +import org.apache.iotdb.commons.exception.runtime.SerializationRunTimeException; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNode; import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId; @@ -41,6 +42,7 @@ import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.exception.NotImplementedException; import org.apache.tsfile.file.metadata.IDeviceID; +import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; import org.apache.tsfile.write.schema.MeasurementSchema; @@ -152,6 +154,12 @@ public void setDataRegionReplicaSet(TRegionReplicaSet dataRegionReplicaSet) { this.dataRegionReplicaSet = dataRegionReplicaSet; } + @Override + public void clearUselessFieldsAfterRouting() { + super.clearUselessFieldsAfterRouting(); + setDataRegionReplicaSet(null); + } + public PartialPath getTargetPath() { return targetPath; } @@ -279,6 +287,43 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { DataNodeQueryMessages.SERIALIZEATTRIBUTES_OF_INSERTNODE_IS_NOT_IMPLEMENTED); } + /** + * Returns the exact number of bytes written by {@link #serializeToByteBuffer()}. + * + * @return the serialized buffer size + */ + public final int serializeToByteBufferSize() { + // InsertNode has no children, so PlanNode.serialize only writes the child count here. + return serializedAttributesSize() + serializedPlanNodeIdSize() + Integer.BYTES; + } + + @Override + public ByteBuffer serializeToByteBuffer() { + try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(serializeToByteBufferSize()); + final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { + serialize(outputStream); + return ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); + } catch (final IOException e) { + throw new SerializationRunTimeException(e); + } + } + + /** + * Returns the exact number of bytes written by the attribute serializer. + * + * @return the serialized attribute size + */ + protected abstract int serializedAttributesSize(); + + /** + * Returns the exact number of bytes written by the plan node id serializer. + * + * @return the serialized plan node id size + */ + protected int serializedPlanNodeIdSize() { + return ReadWriteIOUtils.sizeToWrite(getPlanNodeId().getId()); + } + // region Serialization methods for WAL /** Serialized size of measurement schemas, ignoring failed time series */ diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowNode.java index 04768c57b502f..a5606b4d27a8f 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowNode.java @@ -343,6 +343,73 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { subSerialize(stream); } + @Override + protected int serializedAttributesSize() { + return PlanNodeType.BYTES + serializedSubAttributesSize(); + } + + /** + * Returns the exact number of bytes written by the row serializer. + * + * @return the serialized row field size + */ + protected int serializedSubAttributesSize() { + return Long.BYTES + + ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()) + + serializedMeasurementsAndValuesSize(); + } + + /** + * Returns the exact number of bytes written by the measurement and value serializer. + * + * @return the serialized measurement and value size + */ + protected int serializedMeasurementsAndValuesSize() { + int size = Integer.BYTES + Byte.BYTES; + + for (int i = 0; measurements != null && i < measurements.length; i++) { + if (!shouldSerializeMeasurement(i)) { + continue; + } + size += + measurementSchemas == null + ? ReadWriteIOUtils.sizeToWrite(measurements[i]) + : measurementSchemas[i].serializedSize(); + } + + for (int i = 0; values != null && i < values.length; i++) { + if (!shouldSerializeMeasurement(i)) { + continue; + } + size += serializedValueSize(i); + } + + return size + Byte.BYTES + Byte.BYTES; + } + + private int serializedValueSize(final int index) { + final TSDataType dataType = getDataTypeIfPresent(index); + if (values[index] == null) { + return Byte.BYTES + (dataType == null ? 0 : Byte.BYTES); + } + + if (isNeedInferType) { + return Byte.BYTES + ReadWriteIOUtils.sizeToWrite(values[index].toString()); + } + + return Byte.BYTES + + switch (dataType) { + case BOOLEAN -> Byte.BYTES; + case INT32, DATE -> Integer.BYTES; + case INT64, TIMESTAMP -> Long.BYTES; + case FLOAT -> Float.BYTES; + case DOUBLE -> Double.BYTES; + case TEXT, STRING, BLOB, OBJECT -> ReadWriteIOUtils.sizeToWrite((Binary) values[index]); + case VECTOR, UNKNOWN -> + throw new UnSupportedDataTypeException(UNSUPPORTED_DATA_TYPE + dataType); + }; + } + void subSerialize(ByteBuffer buffer) { ReadWriteIOUtils.write(time, buffer); ReadWriteIOUtils.write(targetPath.getFullPath(), buffer); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java index 4492bf86acde5..7071f9927c08c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java @@ -167,6 +167,12 @@ public void clearResults() { results.clear(); } + @Override + public void clearUselessFieldsAfterRouting() { + super.clearUselessFieldsAfterRouting(); + insertRowNodeList.forEach(InsertRowNode::clearUselessFieldsAfterRouting); + } + public TSStatus[] getFailingStatus() { return StatusUtils.getFailingStatus(results, insertRowNodeList.size()); } @@ -275,6 +281,15 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { } } + @Override + protected int serializedAttributesSize() { + int size = PlanNodeType.BYTES + Integer.BYTES; + for (InsertRowNode node : insertRowNodeList) { + size += node.serializedSubAttributesSize(); + } + return size + insertRowNodeIndexList.size() * Integer.BYTES; + } + @Override public void markAsGeneratedByPipe() { isGeneratedByPipe = true; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java index ccc4ca810d848..cee1e325df228 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java @@ -326,6 +326,16 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { } } + @Override + protected int serializedAttributesSize() { + int size = + PlanNodeType.BYTES + ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()) + Integer.BYTES; + for (InsertRowNode node : insertRowNodeList) { + size += Long.BYTES + node.serializedMeasurementsAndValuesSize(); + } + return size + insertRowNodeIndexList.size() * Integer.BYTES; + } + @Override public void markAsGeneratedByPipe() { isGeneratedByPipe = true; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java index 976223bf2cf10..3f39ee5341203 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java @@ -567,6 +567,89 @@ void subSerialize(DataOutputStream stream) throws IOException { ReadWriteIOUtils.write((byte) (isAligned ? 1 : 0), stream); } + @Override + protected int serializedAttributesSize() { + return PlanNodeType.BYTES + serializedSubAttributesSize(); + } + + /** + * Returns the exact number of bytes written by {@link #subSerialize(DataOutputStream)}. + * + *

This deliberately excludes the plan-node type, id, and children. {@link + * InsertMultiTabletsNode} embeds tablet nodes by calling {@code subSerialize}, rather than their + * complete plan-node serialization. + * + * @return the serialized tablet field size + */ + final int serializedSubAttributesSize() { + int size = ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()); + + size += Integer.BYTES; // valid measurement count + size += Byte.BYTES; // whether measurement schemas are serialized + for (int i = 0; measurements != null && i < measurements.length; i++) { + if (!shouldSerializeMeasurement(i)) { + continue; + } + size += + measurementSchemas == null + ? ReadWriteIOUtils.sizeToWrite(measurements[i]) + : measurementSchemas[i].serializedSize(); + } + + for (int i = 0; dataTypes != null && i < dataTypes.length; i++) { + if (shouldSerializeMeasurement(i)) { + size += TSDataType.getSerializedSize(); + } + } + + size += Integer.BYTES; // row count + size += rowCount * Long.BYTES; // timestamps + + size += Byte.BYTES; // whether bitmaps are serialized + if (bitMaps != null) { + for (int i = 0; measurements != null && i < measurements.length; i++) { + if (!shouldSerializeMeasurement(i)) { + continue; + } + size += Byte.BYTES; // whether the current measurement has a bitmap + if (getBitMapIfPresent(i) != null) { + size += BitMap.getSizeOfBytes(rowCount); + } + } + } + + for (int i = 0; columns != null && i < columns.length; i++) { + if (shouldSerializeMeasurement(i)) { + size += serializedColumnSize(dataTypes[i], columns[i]); + } + } + + return size + Byte.BYTES; // isAligned + } + + private int serializedColumnSize(final TSDataType dataType, final Object column) { + return switch (dataType) { + case BOOLEAN -> rowCount * Byte.BYTES; + case INT32, DATE -> rowCount * Integer.BYTES; + case INT64, TIMESTAMP -> rowCount * Long.BYTES; + case FLOAT -> rowCount * Float.BYTES; + case DOUBLE -> rowCount * Double.BYTES; + case TEXT, BLOB, STRING, OBJECT -> serializedBinaryColumnSize((Binary[]) column); + case VECTOR, UNKNOWN -> + throw new UnSupportedDataTypeException(String.format(DATATYPE_UNSUPPORTED, dataType)); + }; + } + + private int serializedBinaryColumnSize(final Binary[] binaryValues) { + int size = 0; + for (int i = 0; i < rowCount; i++) { + final Binary binary = binaryValues[i]; + final byte[] values = binary == null ? null : binary.getValues(); + size += values == null ? Integer.BYTES : Integer.BYTES + values.length; + } + return size; + } + /** Serialize measurements or measurement schemas, ignoring failed time series */ private void writeMeasurementsOrSchemas(ByteBuffer buffer) { ReadWriteIOUtils.write(getValidMeasurementNumber(), buffer); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertRowNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertRowNode.java index b11c0f6784e3d..468100bf9390b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertRowNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertRowNode.java @@ -237,6 +237,11 @@ void subSerialize(DataOutputStream stream) throws IOException { } } + @Override + protected int serializedSubAttributesSize() { + return super.serializedSubAttributesSize() + getValidMeasurementNumber() * Byte.BYTES; + } + @Override protected void subSerialize(IWALByteBufferView buffer) { super.subSerialize(buffer); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java index d41b078cb9b83..d4c373e57d7f4 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java @@ -324,6 +324,17 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { } } + @Override + protected int serializedAttributesSize() { + int size = super.serializedAttributesSize(); + for (int i = 0; measurements != null && i < measurements.length; i++) { + if (shouldSerializeMeasurement(i)) { + size += Byte.BYTES; + } + } + return size; + } + @Override public void subDeserialize(ByteBuffer buffer) { super.subDeserialize(buffer); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java index df719a724fe4b..7ae9d8d6f77d7 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java @@ -290,6 +290,8 @@ private Future dispatchWrite(List } try { + shouldDispatch.forEach(instance -> instance.getFragment().clearUselessFieldsAfterRouting()); + // 2. try the dispatch final List failedInstances = dispatchWriteOnce(shouldDispatch); 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 6dc5ecb3e2b28..77ff3d8d6e6ed 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 @@ -281,6 +281,21 @@ private void ensureMemTable(long[] infoForMetrics) { } } + private static void clearDataRegionReplicaSet(final InsertRowNode insertRowNode) { + insertRowNode.setDataRegionReplicaSet(null); + } + + private static void clearDataRegionReplicaSet(final InsertRowsNode insertRowsNode) { + insertRowsNode.setDataRegionReplicaSet(null); + for (final InsertRowNode insertRowNode : insertRowsNode.getInsertRowNodeList()) { + clearDataRegionReplicaSet(insertRowNode); + } + } + + private static void clearDataRegionReplicaSet(final InsertTabletNode insertTabletNode) { + insertTabletNode.setDataRegionReplicaSet(null); + } + /** * Insert data in an InsertRowNode into the workingMemtable. * @@ -315,6 +330,7 @@ public void insert(InsertRowNode insertRowNode, long[] infoForMetrics) // recordScheduleMemoryBlockCost infoForMetrics[1] += System.nanoTime() - memControlStartTime; + clearDataRegionReplicaSet(insertRowNode); long startTime = System.nanoTime(); WALFlushListener walFlushListener; try { @@ -414,6 +430,7 @@ public void insertRows(InsertRowsNode insertRowsNode, long[] infoForMetrics) // recordScheduleMemoryBlockCost infoForMetrics[1] += System.nanoTime() - memControlStartTime; + clearDataRegionReplicaSet(insertRowsNode); long startTime = System.nanoTime(); WALFlushListener walFlushListener; try { @@ -586,6 +603,7 @@ public void insertTablet( long[] memIncrements = scheduleMemoryBlock(insertTabletNode, rangeList, results, infoForMetrics); + clearDataRegionReplicaSet(insertTabletNode); long startTime = System.nanoTime(); WALFlushListener walFlushListener; try { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java new file mode 100644 index 0000000000000..e26275fb8e9b2 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java @@ -0,0 +1,530 @@ +/* + * 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.pipe.sink.payload.evolvable.request; + +import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex; +import org.apache.iotdb.commons.exception.IllegalPathException; +import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId; +import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory; +import org.apache.iotdb.db.pipe.sink.protocol.iotconsensusv2.payload.request.IoTConsensusV2TabletInsertNodeReq; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertMultiTabletsNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsOfOneDeviceNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertTabletNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowsNode; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertTabletNode; + +import org.apache.tsfile.enums.ColumnCategory; +import org.apache.tsfile.enums.TSDataType; +import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.apache.tsfile.utils.Binary; +import org.apache.tsfile.utils.BitMap; +import org.apache.tsfile.utils.ReadWriteIOUtils; +import org.apache.tsfile.write.record.Tablet; +import org.apache.tsfile.write.schema.MeasurementSchema; +import org.junit.Assert; +import org.junit.Test; + +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; + +public class PipeTransferSerializationSizeTest { + + @Test + public void testTabletRequestLengths() throws Exception { + final Tablet tablet = createTablet(); + final String database = "pipe_db"; + final PipeTransferTabletRawReq rawReq = + PipeTransferTabletRawReq.toTPipeTransferReq(tablet, false); + assertSerializedBodySize(PipeTransferTabletRawReq.calculateSerializedSize(tablet), rawReq.body); + final PipeTransferTabletRawReqV2 rawReqV2 = + PipeTransferTabletRawReqV2.toTPipeTransferReq(tablet, false, database); + assertSerializedBodySize( + PipeTransferTabletRawReqV2.calculateSerializedSize(tablet, database), rawReqV2.body); + Assert.assertEquals( + PipeTransferTabletRawReq.calculateAirGapSerializedSize(tablet), + PipeTransferTabletRawReq.toTPipeTransferBytes(tablet, false).length); + Assert.assertEquals( + PipeTransferTabletRawReqV2.calculateAirGapSerializedSize(tablet, database), + PipeTransferTabletRawReqV2.toTPipeTransferBytes(tablet, false, database).length); + } + + @Test + public void testBinaryRequestLengths() throws Exception { + final ByteBuffer payload = createByteBufferWithOffsetAndPosition(); + final byte[] expectedPayload = getRemainingBytes(payload); + final String database = "pipe_db_\u6d4b\u8bd5"; + final PipeTransferTabletBinaryReqV2 binaryReqV2 = + PipeTransferTabletBinaryReqV2.toTPipeTransferReq(payload, database); + assertSerializedBodySize( + PipeTransferTabletBinaryReqV2.calculateSerializedSize(payload, database), binaryReqV2.body); + final ByteBuffer thriftBody = binaryReqV2.body.duplicate(); + Assert.assertEquals(expectedPayload.length, ReadWriteIOUtils.readInt(thriftBody)); + assertNextBytes(thriftBody, expectedPayload); + Assert.assertEquals(database, ReadWriteIOUtils.readString(thriftBody)); + Assert.assertFalse(thriftBody.hasRemaining()); + + final byte[] airGapV2Bytes = + PipeTransferTabletBinaryReqV2.toTPipeTransferBytes(payload, database); + Assert.assertEquals( + PipeTransferTabletBinaryReqV2.calculateAirGapSerializedSize(payload, database), + airGapV2Bytes.length); + final ByteBuffer airGapV2Body = ByteBuffer.wrap(airGapV2Bytes); + airGapV2Body.position(Byte.BYTES + Short.BYTES); + Assert.assertEquals(expectedPayload.length, ReadWriteIOUtils.readInt(airGapV2Body)); + assertNextBytes(airGapV2Body, expectedPayload); + Assert.assertEquals(database, ReadWriteIOUtils.readString(airGapV2Body)); + Assert.assertFalse(airGapV2Body.hasRemaining()); + + final byte[] airGapV1Bytes = PipeTransferTabletBinaryReq.toTPipeTransferBytes(payload); + Assert.assertEquals( + PipeTransferTabletBinaryReq.calculateSerializedSize(payload), airGapV1Bytes.length); + final ByteBuffer airGapV1Body = ByteBuffer.wrap(airGapV1Bytes); + airGapV1Body.position(Byte.BYTES + Short.BYTES); + assertNextBytes(airGapV1Body, expectedPayload); + Assert.assertFalse(airGapV1Body.hasRemaining()); + } + + @Test + public void testBatchRequestLengths() throws Exception { + final ByteBuffer insertNode = createByteBufferWithOffsetAndPosition(); + final ByteBuffer tablet = createByteBufferWithOffsetAndPosition(); + final byte[] expectedInsertNode = getRemainingBytes(insertNode); + final byte[] expectedTablet = getRemainingBytes(tablet); + final PipeTransferTabletBatchReq batchReq = + PipeTransferTabletBatchReq.toTPipeTransferReq( + Collections.singletonList(insertNode), Collections.singletonList(tablet)); + assertSerializedBodySize( + PipeTransferTabletBatchReq.calculateSerializedSize( + Collections.singletonList(insertNode), Collections.singletonList(tablet)), + batchReq.body); + final ByteBuffer batchBody = batchReq.body.duplicate(); + Assert.assertEquals(0, ReadWriteIOUtils.readInt(batchBody)); + Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchBody)); + assertNextBytes(batchBody, expectedInsertNode); + Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchBody)); + assertNextBytes(batchBody, expectedTablet); + Assert.assertFalse(batchBody.hasRemaining()); + + final String database = "db_\u6d4b\u8bd5"; + final PipeTransferTabletBatchReqV2 batchReqV2 = + PipeTransferTabletBatchReqV2.toTPipeTransferReq( + Collections.singletonList(insertNode), + Collections.singletonList(tablet), + Collections.singletonList(database), + Collections.singletonList(database)); + assertSerializedBodySize( + PipeTransferTabletBatchReqV2.calculateSerializedSize( + Collections.singletonList(insertNode), + Collections.singletonList(tablet), + Collections.singletonList(database), + Collections.singletonList(database)), + batchReqV2.body); + final ByteBuffer batchV2Body = batchReqV2.body.duplicate(); + Assert.assertEquals(0, ReadWriteIOUtils.readInt(batchV2Body)); + Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchV2Body)); + assertNextBytes(batchV2Body, expectedInsertNode); + Assert.assertEquals(database, ReadWriteIOUtils.readString(batchV2Body)); + Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchV2Body)); + assertNextBytes(batchV2Body, expectedTablet); + Assert.assertEquals(database, ReadWriteIOUtils.readString(batchV2Body)); + Assert.assertFalse(batchV2Body.hasRemaining()); + } + + @Test + public void testInsertNodeSerializedSize() throws Exception { + assertInsertNodeRequestSizes(createInsertRowNode(0), "tree_db"); + assertInsertNodeRequestSizes(createInsertRowNodeWithSchemas(1), "tree_db"); + assertInsertNodeRequestSizes(createInsertRowNodeWithNullValue(), "tree_db"); + assertInsertNodeRequestSizes(createInsertRowNodeWithInferredType(), "tree_db"); + final InsertRowNode partiallyFailedRowNode = createInsertRowNodeWithSchemas(2); + partiallyFailedRowNode.markFailedMeasurement(1); + assertInsertNodeRequestSizes(partiallyFailedRowNode, "tree_db"); + + final InsertTabletNode tabletNode = createInsertTabletNode(); + assertInsertNodeRequestSizes(tabletNode, "tree_db"); + assertInsertNodeRequestSizes(createInsertTabletNode(false, true), "tree_db"); + final InsertTabletNode partiallyFailedTabletNode = createInsertTabletNode(false, true); + partiallyFailedTabletNode.markFailedMeasurement(1); + assertInsertNodeRequestSizes(partiallyFailedTabletNode, "tree_db"); + + final RelationalInsertRowNode relationalRowNode = + new RelationalInsertRowNode( + new PlanNodeId("relational-row"), + new PartialPath("table"), + false, + measurements(), + dataTypes(), + 1, + rowValues(1), + false, + columnCategories()); + assertInsertNodeRequestSizes(relationalRowNode, "table_db_\u6d4b\u8bd5"); + + final RelationalInsertTabletNode relationalTabletNode = createRelationalInsertTabletNode(); + assertInsertNodeRequestSizes(relationalTabletNode, "table_db"); + + final List rows = new ArrayList<>(); + final List relationalRows = new ArrayList<>(); + for (int row = 0; row < 50; row++) { + rows.add(createInsertRowNode(row)); + relationalRows.add(createRelationalInsertRowNode(row)); + } + final InsertRowsNode insertRowsNode = new InsertRowsNode(new PlanNodeId("rows")); + insertRowsNode.setInsertRowNodeList(rows); + insertRowsNode.setInsertRowNodeIndexList(indexes(rows.size())); + assertInsertNodeRequestSizes(insertRowsNode, "tree_db"); + + final InsertRowsOfOneDeviceNode oneDeviceNode = + new InsertRowsOfOneDeviceNode(new PlanNodeId("one-device")); + oneDeviceNode.setInsertRowNodeList(rows); + oneDeviceNode.setInsertRowNodeIndexList(indexes(rows.size())); + assertInsertNodeRequestSizes(oneDeviceNode, "tree_db"); + + final InsertMultiTabletsNode multiTabletsNode = + new InsertMultiTabletsNode(new PlanNodeId("multi-tablets")); + multiTabletsNode.addInsertTabletNode(tabletNode, 0); + multiTabletsNode.addInsertTabletNode(relationalTabletNode, 1); + assertInsertNodeRequestSizes(multiTabletsNode, "tree_db"); + + final RelationalInsertRowsNode relationalRowsNode = + new RelationalInsertRowsNode( + new PlanNodeId("relational-rows"), indexes(relationalRows.size()), relationalRows); + assertInsertNodeRequestSizes(relationalRowsNode, "table_db"); + + final PipeEnrichedInsertNode pipeEnrichedInsertNode = + new PipeEnrichedInsertNode(createInsertRowNode(2)); + pipeEnrichedInsertNode.setPlanNodeId(new PlanNodeId("enriched-row")); + assertInsertNodeRequestSizes(pipeEnrichedInsertNode, "tree_db"); + } + + private static void assertInsertNodeRequestSizes( + final InsertNode insertNode, final String databaseName) throws Exception { + final ByteBuffer serializedInsertNode = insertNode.serializeToByteBuffer(); + Assert.assertEquals(insertNode.serializeToByteBufferSize(), serializedInsertNode.capacity()); + Assert.assertEquals(insertNode.serializeToByteBufferSize(), serializedInsertNode.remaining()); + final PipeTransferTabletInsertNodeReq insertNodeReq = + PipeTransferTabletInsertNodeReq.toTPipeTransferReq(insertNode); + assertSerializedBodySize( + PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode), insertNodeReq.body); + Assert.assertEquals( + PipeTransferTabletInsertNodeReq.calculateAirGapSerializedSize(insertNode), + PipeTransferTabletInsertNodeReq.toTPipeTransferBytes(insertNode).length); + final PipeTransferTabletInsertNodeReqV2 insertNodeReqV2 = + PipeTransferTabletInsertNodeReqV2.toTPipeTransferReq(insertNode, databaseName); + assertSerializedBodySize( + PipeTransferTabletInsertNodeReqV2.calculateSerializedSize(insertNode, databaseName), + insertNodeReqV2.body); + Assert.assertEquals( + PipeTransferTabletInsertNodeReqV2.calculateAirGapSerializedSize(insertNode, databaseName), + PipeTransferTabletInsertNodeReqV2.toTPipeTransferBytes(insertNode, databaseName).length); + final IoTConsensusV2TabletInsertNodeReq iotConsensusReq = + IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq( + insertNode, null, null, MinimumProgressIndex.INSTANCE, 0); + assertSerializedBodySize( + IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode), + iotConsensusReq.body); + } + + private static List indexes(final int size) { + final List indexes = new ArrayList<>(size); + for (int i = 0; i < size; i++) { + indexes.add(i); + } + return indexes; + } + + private static InsertRowNode createInsertRowNode(final int row) throws IllegalPathException { + return new InsertRowNode( + new PlanNodeId("row-" + row), + new PartialPath("root.sg.d"), + false, + measurements(), + dataTypes(), + row, + rowValues(row), + false); + } + + private static InsertRowNode createInsertRowNodeWithSchemas(final int row) + throws IllegalPathException { + return new InsertRowNode( + new PlanNodeId("row-with-schemas"), + new PartialPath("root.sg.d"), + false, + measurements(), + dataTypes(), + measurementSchemas(), + row, + rowValues(row), + false); + } + + private static InsertRowNode createInsertRowNodeWithNullValue() throws IllegalPathException { + return new InsertRowNode( + new PlanNodeId("row-with-null"), + new PartialPath("root.sg.d"), + false, + new String[] {"s"}, + new TSDataType[] {TSDataType.INT32}, + 1, + new Object[] {null}, + false); + } + + private static InsertRowNode createInsertRowNodeWithInferredType() throws IllegalPathException { + return new InsertRowNode( + new PlanNodeId("row-with-inferred-type"), + new PartialPath("root.sg.d"), + false, + new String[] {"s"}, + new TSDataType[] {null}, + 1, + new Object[] {"value"}, + true); + } + + private static RelationalInsertRowNode createRelationalInsertRowNode(final int row) + throws IllegalPathException { + return new RelationalInsertRowNode( + new PlanNodeId("relational-row-" + row), + new PartialPath("table"), + false, + measurements(), + dataTypes(), + row, + rowValues(row), + false, + columnCategories()); + } + + private static InsertTabletNode createInsertTabletNode() throws IllegalPathException { + return createInsertTabletNode(false, false); + } + + private static InsertTabletNode createInsertTabletNode(final boolean relational) + throws IllegalPathException { + return createInsertTabletNode(relational, false); + } + + private static InsertTabletNode createInsertTabletNode( + final boolean relational, final boolean withBitMaps) throws IllegalPathException { + final String[] measurements = measurements(); + final TSDataType[] types = dataTypes(); + final MeasurementSchema[] schemas = measurementSchemas(); + final Object[] columns = new Object[types.length]; + final int rowCount = 50; + final long[] times = new long[rowCount]; + for (int i = 0; i < rowCount; i++) { + times[i] = i; + } + for (int column = 0; column < types.length; column++) { + switch (types[column]) { + case BOOLEAN: + final boolean[] booleanValues = new boolean[rowCount]; + for (int row = 0; row < rowCount; row++) { + booleanValues[row] = row % 2 == 0; + } + columns[column] = booleanValues; + break; + case INT32: + case DATE: + final int[] intValues = new int[rowCount]; + for (int row = 0; row < rowCount; row++) { + intValues[row] = row; + } + columns[column] = intValues; + break; + case INT64: + case TIMESTAMP: + final long[] longValues = new long[rowCount]; + for (int row = 0; row < rowCount; row++) { + longValues[row] = row; + } + columns[column] = longValues; + break; + case FLOAT: + final float[] floatValues = new float[rowCount]; + for (int row = 0; row < rowCount; row++) { + floatValues[row] = row; + } + columns[column] = floatValues; + break; + case DOUBLE: + final double[] doubleValues = new double[rowCount]; + for (int row = 0; row < rowCount; row++) { + doubleValues[row] = row; + } + columns[column] = doubleValues; + break; + case TEXT: + case BLOB: + case STRING: + case OBJECT: + Binary[] values = new Binary[rowCount]; + for (int row = 1; row < rowCount; row++) { + values[row] = new Binary(("value-" + row).getBytes(StandardCharsets.UTF_8)); + } + columns[column] = values; + break; + default: + throw new AssertionError(types[column]); + } + } + final BitMap[] bitMaps = withBitMaps ? createBitMaps(types.length, rowCount) : null; + return relational + ? new RelationalInsertTabletNode( + new PlanNodeId("relational-tablet"), + new PartialPath("table"), + false, + measurements, + types, + schemas, + times, + bitMaps, + columns, + rowCount, + columnCategories()) + : new InsertTabletNode( + new PlanNodeId("tablet"), + new PartialPath("root.sg.d"), + false, + measurements, + types, + schemas, + times, + bitMaps, + columns, + rowCount); + } + + private static RelationalInsertTabletNode createRelationalInsertTabletNode() + throws IllegalPathException { + return (RelationalInsertTabletNode) createInsertTabletNode(true); + } + + private static String[] measurements() { + return new String[] {"b", "i", "l", "f", "d", "t", "ts", "date", "blob", "string", "object"}; + } + + private static TSDataType[] dataTypes() { + return new TSDataType[] { + TSDataType.BOOLEAN, + TSDataType.INT32, + TSDataType.INT64, + TSDataType.FLOAT, + TSDataType.DOUBLE, + TSDataType.TEXT, + TSDataType.TIMESTAMP, + TSDataType.DATE, + TSDataType.BLOB, + TSDataType.STRING, + TSDataType.OBJECT + }; + } + + private static MeasurementSchema[] measurementSchemas() { + final String[] measurements = measurements(); + final TSDataType[] types = dataTypes(); + final MeasurementSchema[] schemas = new MeasurementSchema[types.length]; + for (int i = 0; i < types.length; i++) { + schemas[i] = new MeasurementSchema(measurements[i], types[i], TSEncoding.PLAIN); + } + return schemas; + } + + private static BitMap[] createBitMaps(final int columnCount, final int rowCount) { + final BitMap[] bitMaps = new BitMap[columnCount]; + bitMaps[0] = new BitMap(rowCount); + bitMaps[0].mark(0); + bitMaps[columnCount - 1] = new BitMap(rowCount); + bitMaps[columnCount - 1].mark(rowCount - 1); + return bitMaps; + } + + private static Object[] rowValues(final int row) { + return new Object[] { + true, + row, + (long) row, + (float) row, + (double) row, + new Binary(("text-" + row).getBytes(StandardCharsets.UTF_8)), + (long) row, + row, + new Binary(("blob-" + row).getBytes(StandardCharsets.UTF_8)), + new Binary(("string-" + row).getBytes(StandardCharsets.UTF_8)), + new Binary(("object-" + row).getBytes(StandardCharsets.UTF_8)) + }; + } + + private static TsTableColumnCategory[] columnCategories() { + final TsTableColumnCategory[] categories = new TsTableColumnCategory[dataTypes().length]; + Arrays.fill(categories, TsTableColumnCategory.FIELD); + categories[0] = TsTableColumnCategory.TAG; + return categories; + } + + private static ByteBuffer createByteBufferWithOffsetAndPosition() { + final ByteBuffer source = ByteBuffer.wrap(new byte[] {0, 1, 2, 3, 4, 5}); + source.position(1); + final ByteBuffer buffer = source.slice(); + buffer.position(1); + buffer.limit(4); + return buffer; + } + + private static byte[] getRemainingBytes(final ByteBuffer buffer) { + final ByteBuffer duplicate = buffer.duplicate(); + final byte[] bytes = new byte[duplicate.remaining()]; + duplicate.get(bytes); + return bytes; + } + + private static void assertSerializedBodySize(final int expectedSize, final ByteBuffer body) { + Assert.assertEquals(expectedSize, body.remaining()); + Assert.assertEquals(expectedSize, body.capacity()); + } + + private static void assertNextBytes(final ByteBuffer buffer, final byte[] expectedBytes) { + final byte[] actualBytes = new byte[expectedBytes.length]; + buffer.get(actualBytes); + Assert.assertArrayEquals(expectedBytes, actualBytes); + } + + private static Tablet createTablet() { + final Tablet tablet = + new Tablet( + "table1", Collections.singletonList(new MeasurementSchema("s1", TSDataType.INT32)), 1); + tablet.setColumnCategories(Collections.singletonList(ColumnCategory.FIELD)); + tablet.addTimestamp(0, 1L); + tablet.addValue(0, 0, 1); + tablet.setRowSize(1); + return tablet; + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java index 907af18ef9adb..e09a71e435137 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java @@ -19,11 +19,13 @@ package org.apache.iotdb.db.queryengine.plan.planner.node.write; +import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeType; import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory; +import org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowNode; @@ -45,6 +47,30 @@ public class InsertRowsNodeSerdeTest { + @Test + public void testClearUselessFieldAfterDispatch() throws IllegalPathException { + final InsertRowsNode insertRowsNode = new InsertRowsNode(new PlanNodeId("insert rows")); + final InsertRowNode insertRowNode = + new InsertRowNode( + new PlanNodeId("insert row"), + new PartialPath("root.sg.d1"), + false, + new String[] {"s1"}, + new TSDataType[] {TSDataType.INT32}, + 1L, + new Object[] {1}, + false); + insertRowsNode.addOneInsertRowNode(insertRowNode, 0); + + insertRowsNode.setDataRegionReplicaSet(new TRegionReplicaSet()); + insertRowNode.setDataRegionReplicaSet(new TRegionReplicaSet()); + + new PipeEnrichedInsertNode(insertRowsNode).clearUselessFieldsAfterRouting(); + + Assert.assertNull(insertRowsNode.getDataRegionReplicaSet()); + Assert.assertNull(insertRowNode.getDataRegionReplicaSet()); + } + @Test public void TestSerializeAndDeserialize() throws IllegalPathException { InsertRowsNode node = new InsertRowsNode(new PlanNodeId("plan node 1")); diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java index b2e74b32af847..3411efde47e33 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java @@ -78,6 +78,14 @@ public void markAsGeneratedByPipe() { public abstract void addChild(PlanNode child); + /** Releases fields that are no longer needed after the target region has been determined. */ + public void clearUselessFieldsAfterRouting() { + final List children = getChildren(); + if (children != null) { + children.forEach(PlanNode::clearUselessFieldsAfterRouting); + } + } + /** * If this plan node has to be serialized or deserialized, override this method. If this method is * overridden, the serialization and deserialization methods must be implemented.