From 793b0a6474bd29491133587e1fb558642f693a4f Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Fri, 17 Jul 2026 14:26:42 +0800 Subject: [PATCH 1/5] Optimize pipe request serialization buffer sizing --- .../batch/PipeTabletEventPlainBatch.java | 13 ++- .../request/PipeTransferTabletBatchReq.java | 10 +- .../request/PipeTransferTabletBatchReqV2.java | 27 ++++- .../request/PipeTransferTabletBinaryReq.java | 6 +- .../PipeTransferTabletBinaryReqV2.java | 19 +++- .../request/PipeTransferTabletRawReq.java | 13 ++- .../request/PipeTransferTabletRawReqV2.java | 16 ++- .../PipeTransferSerializationSizeTest.java | 101 ++++++++++++++++++ 8 files changed, 193 insertions(+), 12 deletions(-) create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java 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..832287c7148ea 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); @@ -192,9 +193,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()); } @@ -234,6 +237,10 @@ 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..917547b278e5c 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 @@ -129,7 +129,8 @@ 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); @@ -151,6 +152,13 @@ public static PipeTransferTabletBatchReq toTPipeTransferReq( return batchReq; } + static int calculateSerializedSize( + final List insertNodeBuffers, final List tabletBuffers) { + return Integer.BYTES * 3 + + insertNodeBuffers.stream().mapToInt(ByteBuffer::limit).sum() + + tabletBuffers.stream().mapToInt(ByteBuffer::limit).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..2415aaf16c5d5 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 @@ -33,6 +33,7 @@ import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement; import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq; +import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; @@ -186,7 +187,10 @@ 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); @@ -212,6 +216,27 @@ public static PipeTransferTabletBatchReqV2 toTPipeTransferReq( return batchReq; } + static int calculateSerializedSize( + final List insertNodeBuffers, + final List tabletBuffers, + final List insertNodeDataBases, + final List tabletDataBases) { + int size = Integer.BYTES * 3; + for (int i = 0; i < insertNodeBuffers.size(); i++) { + size += insertNodeBuffers.get(i).limit(); + size += serializedStringSize(insertNodeDataBases.get(i)); + } + for (int i = 0; i < tabletBuffers.size(); i++) { + size += tabletBuffers.get(i).limit(); + size += serializedStringSize(tabletDataBases.get(i)); + } + return size; + } + + private static int serializedStringSize(final String value) { + return Integer.BYTES + (value == null ? 0 : value.getBytes(TSFileConfig.STRING_CHARSET).length); + } + 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..0d790ff1bc1ba 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 @@ -102,7 +102,7 @@ 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(Byte.BYTES + Short.BYTES); final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), outputStream); ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_BINARY.getType(), outputStream); @@ -110,6 +110,10 @@ public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer) throws IO } } + static int calculateSerializedSize(final ByteBuffer byteBuffer) { + return Byte.BYTES + Short.BYTES + byteBuffer.limit(); + } + /////////////////////////////// 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..f80925da14411 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 @@ -32,6 +32,7 @@ import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement; import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement; +import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; @@ -119,7 +120,8 @@ 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()); @@ -151,7 +153,8 @@ 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); @@ -162,6 +165,18 @@ public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer, final Str } } + private static int serializedStringSize(final String value) { + return Integer.BYTES + (value == null ? 0 : value.getBytes(TSFileConfig.STRING_CHARSET).length); + } + + static int calculateSerializedSize(final ByteBuffer byteBuffer, final String dataBaseName) { + return Integer.BYTES + byteBuffer.limit() + serializedStringSize(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/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..1b6c24d571663 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,9 @@ 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( + Byte.BYTES + Short.BYTES + calculateSerializedSize(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 +187,15 @@ public static byte[] toTPipeTransferBytes( } } + static int calculateSerializedSize(final Tablet tablet, final String dataBaseName) { + return tablet.serializedSize() + + Byte.BYTES + + Integer.BYTES + + (dataBaseName == null + ? 0 + : dataBaseName.getBytes(java.nio.charset.StandardCharsets.UTF_8).length); + } + /////////////////////////////// Object /////////////////////////////// @Override 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..3890ced89897d --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java @@ -0,0 +1,101 @@ +/* + * 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.tsfile.enums.ColumnCategory; +import org.apache.tsfile.enums.TSDataType; +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.util.Collections; + +public class PipeTransferSerializationSizeTest { + + @Test + public void testTabletRequestLengths() throws Exception { + final Tablet tablet = createTablet(); + final String database = "pipe_db"; + Assert.assertEquals( + PipeTransferTabletRawReq.calculateSerializedSize(tablet), + PipeTransferTabletRawReq.toTPipeTransferReq(tablet, false).getBody().length); + Assert.assertEquals( + PipeTransferTabletRawReqV2.calculateSerializedSize(tablet, database), + PipeTransferTabletRawReqV2.toTPipeTransferReq(tablet, false, database).getBody().length); + Assert.assertEquals( + PipeTransferTabletRawReq.calculateAirGapSerializedSize(tablet), + PipeTransferTabletRawReq.toTPipeTransferBytes(tablet, false).length); + } + + @Test + public void testBinaryRequestLengths() throws Exception { + final ByteBuffer payload = ByteBuffer.wrap(new byte[] {1, 2, 3, 4}); + final String database = "pipe_db"; + Assert.assertEquals( + PipeTransferTabletBinaryReqV2.calculateSerializedSize(payload, database), + PipeTransferTabletBinaryReqV2.toTPipeTransferReq(payload, database).getBody().length); + Assert.assertEquals( + PipeTransferTabletBinaryReqV2.calculateAirGapSerializedSize(payload, database), + PipeTransferTabletBinaryReqV2.toTPipeTransferBytes(payload, database).length); + Assert.assertEquals( + PipeTransferTabletBinaryReq.calculateSerializedSize(payload), + PipeTransferTabletBinaryReq.toTPipeTransferBytes(payload).length); + } + + @Test + public void testBatchRequestLengths() throws Exception { + final ByteBuffer insertNode = ByteBuffer.wrap(new byte[] {1, 2}); + final ByteBuffer tablet = ByteBuffer.wrap(new byte[] {3, 4, 5}); + Assert.assertEquals( + PipeTransferTabletBatchReq.calculateSerializedSize( + Collections.singletonList(insertNode), Collections.singletonList(tablet)), + PipeTransferTabletBatchReq.toTPipeTransferReq( + Collections.singletonList(insertNode), Collections.singletonList(tablet)) + .getBody() + .length); + + final String database = "db"; + Assert.assertEquals( + PipeTransferTabletBatchReqV2.calculateSerializedSize( + Collections.singletonList(insertNode), + Collections.singletonList(tablet), + Collections.singletonList(database), + Collections.singletonList(database)), + PipeTransferTabletBatchReqV2.toTPipeTransferReq( + Collections.singletonList(insertNode), + Collections.singletonList(tablet), + Collections.singletonList(database), + Collections.singletonList(database)) + .getBody() + .length); + } + + 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; + } +} From e6be569a47529e816437c5c4fd4b3edf10942b93 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Mon, 20 Jul 2026 14:11:47 +0800 Subject: [PATCH 2/5] Add exact insert node serialization sizing --- .../PipeTransferTabletInsertNodeReq.java | 30 +- .../PipeTransferTabletInsertNodeReqV2.java | 16 +- ...IoTConsensusV2TransferBatchReqBuilder.java | 9 +- .../IoTConsensusV2TabletInsertNodeReq.java | 6 + .../node/pipe/PipeEnrichedInsertNode.java | 11 + .../node/write/InsertMultiTabletsNode.java | 9 + .../planner/plan/node/write/InsertNode.java | 27 ++ .../plan/node/write/InsertRowNode.java | 69 ++++ .../plan/node/write/InsertRowsNode.java | 9 + .../node/write/InsertRowsOfOneDeviceNode.java | 10 + .../plan/node/write/InsertTabletNode.java | 91 +++++ .../node/write/RelationalInsertRowNode.java | 5 + .../write/RelationalInsertTabletNode.java | 11 + .../PipeTransferSerializationSizeTest.java | 357 ++++++++++++++++++ 14 files changed, 647 insertions(+), 13 deletions(-) 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..14033fb5f1354 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,12 +31,12 @@ 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; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Objects; public class PipeTransferTabletInsertNodeReq extends TPipeTransferReq { @@ -86,7 +86,14 @@ public static PipeTransferTabletInsertNodeReq toTPipeTransferReq(final InsertNod req.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); req.type = PipeRequestType.TRANSFER_TABLET_INSERT_NODE.getType(); - req.body = insertNode.serializeToByteBuffer(); + try (final PublicBAOS byteArrayOutputStream = + new PublicBAOS(calculateSerializedSize(insertNode)); + final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { + insertNode.serialize(outputStream); + req.body = ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); + } catch (final IOException e) { + throw new RuntimeException(e); + } return req; } @@ -106,15 +113,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/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..8b83c9aba6957 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,16 @@ 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(); + final int serializedSize = + IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode); batchReqs.add( IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq( insertNode, commitId, consensusGroupId, progressIndex, thisDataNodeId)); - return buffer.limit(); + return serializedSize; } @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/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..2abff39a65ce9 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; @@ -290,6 +291,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..70f5435f6d789 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.baseSubSerializedSizeForPipe(); + } + 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..6de154224f439 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; @@ -279,6 +281,31 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { DataNodeQueryMessages.SERIALIZEATTRIBUTES_OF_INSERTNODE_IS_NOT_IMPLEMENTED); } + /** Returns the exact size of the buffer produced by {@link #serializeToByteBuffer()}. */ + 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 size of the attributes written by {@link #serializeToByteBuffer()}. */ + protected abstract int serializedAttributesSize(); + + /** Returns the exact size of the plan node id written after the attributes. */ + 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..67671865951bc 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,75 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { subSerialize(stream); } + @Override + protected int serializedAttributesSize() { + return PlanNodeType.BYTES + pipeSubSerializedSize(); + } + + /** Returns the exact size of the row fields written during Pipe serialization. */ + protected int pipeSubSerializedSize() { + return Long.BYTES + + ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()) + + pipeMeasurementsAndValuesSerializedSize(); + } + + /** Returns the exact size of measurement and value fields written during Pipe serialization. */ + protected int pipeMeasurementsAndValuesSerializedSize() { + 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 += pipeValueSerializedSize(i); + } + + return size + Byte.BYTES + Byte.BYTES; + } + + private int pipeValueSerializedSize(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()); + } + + switch (dataType) { + case BOOLEAN: + return Byte.BYTES + Byte.BYTES; + case INT32: + case DATE: + return Byte.BYTES + Integer.BYTES; + case INT64: + case TIMESTAMP: + return Byte.BYTES + Long.BYTES; + case FLOAT: + return Byte.BYTES + Float.BYTES; + case DOUBLE: + return Byte.BYTES + Double.BYTES; + case TEXT: + case STRING: + case BLOB: + case OBJECT: + return Byte.BYTES + ReadWriteIOUtils.sizeToWrite((Binary) values[index]); + default: + 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..ec46e4104aaf2 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 @@ -275,6 +275,15 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { } } + @Override + protected int serializedAttributesSize() { + int size = PlanNodeType.BYTES + Integer.BYTES; + for (InsertRowNode node : insertRowNodeList) { + size += node.pipeSubSerializedSize(); + } + 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..17eb73a3e2947 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.pipeMeasurementsAndValuesSerializedSize(); + } + 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..c8a979f7e2f11 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,97 @@ void subSerialize(DataOutputStream stream) throws IOException { ReadWriteIOUtils.write((byte) (isAligned ? 1 : 0), stream); } + @Override + protected int serializedAttributesSize() { + return PlanNodeType.BYTES + baseSubSerializedSizeForPipe(); + } + + /** + * Returns the exact size 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. + */ + final int baseSubSerializedSizeForPipe() { + int size = ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()); + + size += Integer.BYTES; + size += 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; dataTypes != null && i < dataTypes.length; i++) { + if (shouldSerializeMeasurement(i)) { + size += TSDataType.getSerializedSize(); + } + } + + size += Integer.BYTES; + size += rowCount * Long.BYTES; + + size += Byte.BYTES; + if (bitMaps != null) { + for (int i = 0; measurements != null && i < measurements.length; i++) { + if (!shouldSerializeMeasurement(i)) { + continue; + } + size += Byte.BYTES; + if (getBitMapIfPresent(i) != null) { + size += BitMap.getSizeOfBytes(rowCount); + } + } + } + + for (int i = 0; columns != null && i < columns.length; i++) { + if (shouldSerializeMeasurement(i)) { + size += columnSerializedSizeForPipe(dataTypes[i], columns[i]); + } + } + + return size + Byte.BYTES; + } + + private int columnSerializedSizeForPipe(final TSDataType dataType, final Object column) { + switch (dataType) { + case INT32: + case DATE: + return rowCount * Integer.BYTES; + case INT64: + case TIMESTAMP: + return rowCount * Long.BYTES; + case FLOAT: + return rowCount * Float.BYTES; + case DOUBLE: + return rowCount * Double.BYTES; + case BOOLEAN: + return rowCount * Byte.BYTES; + case TEXT: + case BLOB: + case STRING: + case OBJECT: + int size = 0; + final Binary[] binaryValues = (Binary[]) column; + for (int i = 0; i < rowCount; i++) { + final Binary binary = binaryValues[i]; + size += + binary == null || binary.getValues() == null + ? Integer.BYTES + : Integer.BYTES + binary.getValues().length; + } + return size; + default: + throw new UnSupportedDataTypeException(String.format(DATATYPE_UNSUPPORTED, dataType)); + } + } + /** 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..ae5abf83377cc 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 pipeSubSerializedSize() { + return super.pipeSubSerializedSize() + 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/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 index 3890ced89897d..ebaafe4fbcd2c 100644 --- 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 @@ -18,15 +18,39 @@ 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.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 { @@ -88,6 +112,339 @@ public void testBatchRequestLengths() throws Exception { .length); } + @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()); + Assert.assertEquals( + PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode), + PipeTransferTabletInsertNodeReq.toTPipeTransferReq(insertNode).getBody().length); + Assert.assertEquals( + PipeTransferTabletInsertNodeReq.calculateAirGapSerializedSize(insertNode), + PipeTransferTabletInsertNodeReq.toTPipeTransferBytes(insertNode).length); + Assert.assertEquals( + PipeTransferTabletInsertNodeReqV2.calculateSerializedSize(insertNode, databaseName), + PipeTransferTabletInsertNodeReqV2.toTPipeTransferReq(insertNode, databaseName) + .getBody() + .length); + Assert.assertEquals( + PipeTransferTabletInsertNodeReqV2.calculateAirGapSerializedSize(insertNode, databaseName), + PipeTransferTabletInsertNodeReqV2.toTPipeTransferBytes(insertNode, databaseName).length); + Assert.assertEquals( + IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode), + IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq( + insertNode, null, null, MinimumProgressIndex.INSTANCE, 0) + .getBody() + .length); + } + + 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 Tablet createTablet() { final Tablet tablet = new Tablet( From 66a111e7b6b6322c2f2df0f48ed08164b340c971 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Mon, 20 Jul 2026 18:57:08 +0800 Subject: [PATCH 3/5] Optimize Pipe serialization buffer sizing --- .../request/PipeTransferTabletBatchReq.java | 23 ++- .../request/PipeTransferTabletBatchReqV2.java | 33 ++-- .../request/PipeTransferTabletBinaryReq.java | 12 +- .../PipeTransferTabletBinaryReqV2.java | 21 +-- .../PipeTransferTabletInsertNodeReq.java | 10 +- .../request/PipeTransferTabletRawReqV2.java | 14 +- ...IoTConsensusV2TransferBatchReqBuilder.java | 9 +- .../node/write/InsertMultiTabletsNode.java | 2 +- .../planner/plan/node/write/InsertNode.java | 18 +- .../plan/node/write/InsertRowNode.java | 56 +++---- .../plan/node/write/InsertRowsNode.java | 2 +- .../node/write/InsertRowsOfOneDeviceNode.java | 2 +- .../plan/node/write/InsertTabletNode.java | 74 ++++---- .../node/write/RelationalInsertRowNode.java | 4 +- .../PipeTransferSerializationSizeTest.java | 158 +++++++++++++----- 15 files changed, 264 insertions(+), 174 deletions(-) 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 917547b278e5c..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<>(); @@ -135,14 +140,22 @@ public static PipeTransferTabletBatchReq toTPipeTransferReq( // 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 = @@ -154,9 +167,9 @@ public static PipeTransferTabletBatchReq toTPipeTransferReq( static int calculateSerializedSize( final List insertNodeBuffers, final List tabletBuffers) { - return Integer.BYTES * 3 - + insertNodeBuffers.stream().mapToInt(ByteBuffer::limit).sum() - + tabletBuffers.stream().mapToInt(ByteBuffer::limit).sum(); + return BATCH_REQUEST_COUNT_SERIALIZED_SIZE + + insertNodeBuffers.stream().mapToInt(ByteBuffer::remaining).sum() + + tabletBuffers.stream().mapToInt(ByteBuffer::remaining).sum(); } public static PipeTransferTabletBatchReq fromTPipeTransferReq( 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 2415aaf16c5d5..c1d388031bb4f 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 @@ -33,7 +33,6 @@ import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement; import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq; -import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; @@ -47,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<>(); @@ -195,17 +200,25 @@ public static PipeTransferTabletBatchReqV2 toTPipeTransferReq( // 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); } @@ -221,22 +234,18 @@ static int calculateSerializedSize( final List tabletBuffers, final List insertNodeDataBases, final List tabletDataBases) { - int size = Integer.BYTES * 3; + int size = BATCH_REQUEST_COUNT_SERIALIZED_SIZE; for (int i = 0; i < insertNodeBuffers.size(); i++) { - size += insertNodeBuffers.get(i).limit(); - size += serializedStringSize(insertNodeDataBases.get(i)); + size += insertNodeBuffers.get(i).remaining(); + size += ReadWriteIOUtils.sizeToWrite(insertNodeDataBases.get(i)); } for (int i = 0; i < tabletBuffers.size(); i++) { - size += tabletBuffers.get(i).limit(); - size += serializedStringSize(tabletDataBases.get(i)); + size += tabletBuffers.get(i).remaining(); + size += ReadWriteIOUtils.sizeToWrite(tabletDataBases.get(i)); } return size; } - private static int serializedStringSize(final String value) { - return Integer.BYTES + (value == null ? 0 : value.getBytes(TSFileConfig.STRING_CHARSET).length); - } - 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 0d790ff1bc1ba..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,16 +101,21 @@ public static PipeTransferTabletBinaryReq fromTPipeTransferReq( /////////////////////////////// Air Gap /////////////////////////////// public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer) throws IOException { - try (final PublicBAOS byteArrayOutputStream = new PublicBAOS(Byte.BYTES + Short.BYTES); + 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.limit(); + return Byte.BYTES + Short.BYTES + byteBuffer.remaining(); } /////////////////////////////// Object /////////////////////////////// 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 f80925da14411..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 @@ -32,7 +32,6 @@ import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement; import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement; -import org.apache.tsfile.common.conf.TSFileConfig; import org.apache.tsfile.utils.PublicBAOS; import org.apache.tsfile.utils.ReadWriteIOUtils; @@ -123,8 +122,11 @@ public static PipeTransferTabletBinaryReqV2 toTPipeTransferReq( 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()); } @@ -158,19 +160,18 @@ public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer, final Str 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(); } } - private static int serializedStringSize(final String value) { - return Integer.BYTES + (value == null ? 0 : value.getBytes(TSFileConfig.STRING_CHARSET).length); - } - static int calculateSerializedSize(final ByteBuffer byteBuffer, final String dataBaseName) { - return Integer.BYTES + byteBuffer.limit() + serializedStringSize(dataBaseName); + return Integer.BYTES + byteBuffer.remaining() + ReadWriteIOUtils.sizeToWrite(dataBaseName); } static int calculateAirGapSerializedSize(final ByteBuffer byteBuffer, final String dataBaseName) { 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 14033fb5f1354..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 @@ -36,7 +36,6 @@ import java.io.DataOutputStream; import java.io.IOException; -import java.nio.ByteBuffer; import java.util.Objects; public class PipeTransferTabletInsertNodeReq extends TPipeTransferReq { @@ -86,14 +85,7 @@ public static PipeTransferTabletInsertNodeReq toTPipeTransferReq(final InsertNod req.version = IoTDBSinkRequestVersion.VERSION_1.getVersion(); req.type = PipeRequestType.TRANSFER_TABLET_INSERT_NODE.getType(); - try (final PublicBAOS byteArrayOutputStream = - new PublicBAOS(calculateSerializedSize(insertNode)); - final DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { - insertNode.serialize(outputStream); - req.body = ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); - } catch (final IOException e) { - throw new RuntimeException(e); - } + req.body = insertNode.serializeToByteBuffer(); return req; } 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 1b6c24d571663..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 @@ -175,8 +175,7 @@ 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( - Byte.BYTES + Short.BYTES + calculateSerializedSize(tablet, dataBaseName)); + 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); @@ -188,12 +187,11 @@ public static byte[] toTPipeTransferBytes( } static int calculateSerializedSize(final Tablet tablet, final String dataBaseName) { - return tablet.serializedSize() - + Byte.BYTES - + Integer.BYTES - + (dataBaseName == null - ? 0 - : dataBaseName.getBytes(java.nio.charset.StandardCharsets.UTF_8).length); + 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 /////////////////////////////// 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 8b83c9aba6957..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 @@ -222,13 +222,12 @@ protected int buildTabletInsertionBuffer(TabletInsertionEvent event) throws WALP final InsertNode insertNode = pipeInsertNodeTabletInsertionEvent.getInsertNode(); // IoTConsensusV2 will transfer binary data to TIoTConsensusV2TransferReq final ProgressIndex progressIndex = pipeInsertNodeTabletInsertionEvent.getProgressIndex(); - final int serializedSize = - IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode); - batchReqs.add( + final IoTConsensusV2TabletInsertNodeReq request = IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq( - insertNode, commitId, consensusGroupId, progressIndex, thisDataNodeId)); + insertNode, commitId, consensusGroupId, progressIndex, thisDataNodeId); + batchReqs.add(request); - return serializedSize; + return request.body.remaining(); } @Override 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 70f5435f6d789..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 @@ -284,7 +284,7 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { protected int serializedAttributesSize() { int size = PlanNodeType.BYTES + Integer.BYTES; for (final InsertTabletNode insertTabletNode : insertTabletNodeList) { - size += insertTabletNode.baseSubSerializedSizeForPipe(); + size += insertTabletNode.serializedSubAttributesSize(); } return size + parentInsertTabletNodeIndexList.size() * Integer.BYTES; } 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 6de154224f439..37c1cdb734f39 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 @@ -281,7 +281,11 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { DataNodeQueryMessages.SERIALIZEATTRIBUTES_OF_INSERTNODE_IS_NOT_IMPLEMENTED); } - /** Returns the exact size of the buffer produced by {@link #serializeToByteBuffer()}. */ + /** + * 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; @@ -298,10 +302,18 @@ public ByteBuffer serializeToByteBuffer() { } } - /** Returns the exact size of the attributes written by {@link #serializeToByteBuffer()}. */ + /** + * Returns the exact number of bytes written by the attribute serializer. + * + * @return the serialized attribute size + */ protected abstract int serializedAttributesSize(); - /** Returns the exact size of the plan node id written after the attributes. */ + /** + * 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()); } 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 67671865951bc..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 @@ -345,18 +345,26 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { @Override protected int serializedAttributesSize() { - return PlanNodeType.BYTES + pipeSubSerializedSize(); + return PlanNodeType.BYTES + serializedSubAttributesSize(); } - /** Returns the exact size of the row fields written during Pipe serialization. */ - protected int pipeSubSerializedSize() { + /** + * 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()) - + pipeMeasurementsAndValuesSerializedSize(); + + serializedMeasurementsAndValuesSize(); } - /** Returns the exact size of measurement and value fields written during Pipe serialization. */ - protected int pipeMeasurementsAndValuesSerializedSize() { + /** + * 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++) { @@ -373,13 +381,13 @@ protected int pipeMeasurementsAndValuesSerializedSize() { if (!shouldSerializeMeasurement(i)) { continue; } - size += pipeValueSerializedSize(i); + size += serializedValueSize(i); } return size + Byte.BYTES + Byte.BYTES; } - private int pipeValueSerializedSize(final int index) { + private int serializedValueSize(final int index) { final TSDataType dataType = getDataTypeIfPresent(index); if (values[index] == null) { return Byte.BYTES + (dataType == null ? 0 : Byte.BYTES); @@ -389,27 +397,17 @@ private int pipeValueSerializedSize(final int index) { return Byte.BYTES + ReadWriteIOUtils.sizeToWrite(values[index].toString()); } - switch (dataType) { - case BOOLEAN: - return Byte.BYTES + Byte.BYTES; - case INT32: - case DATE: - return Byte.BYTES + Integer.BYTES; - case INT64: - case TIMESTAMP: - return Byte.BYTES + Long.BYTES; - case FLOAT: - return Byte.BYTES + Float.BYTES; - case DOUBLE: - return Byte.BYTES + Double.BYTES; - case TEXT: - case STRING: - case BLOB: - case OBJECT: - return Byte.BYTES + ReadWriteIOUtils.sizeToWrite((Binary) values[index]); - default: - throw new UnSupportedDataTypeException(UNSUPPORTED_DATA_TYPE + dataType); - } + 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) { 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 ec46e4104aaf2..fe36020a55e46 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 @@ -279,7 +279,7 @@ protected void serializeAttributes(DataOutputStream stream) throws IOException { protected int serializedAttributesSize() { int size = PlanNodeType.BYTES + Integer.BYTES; for (InsertRowNode node : insertRowNodeList) { - size += node.pipeSubSerializedSize(); + size += node.serializedSubAttributesSize(); } return size + insertRowNodeIndexList.size() * Integer.BYTES; } 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 17eb73a3e2947..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 @@ -331,7 +331,7 @@ protected int serializedAttributesSize() { int size = PlanNodeType.BYTES + ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()) + Integer.BYTES; for (InsertRowNode node : insertRowNodeList) { - size += Long.BYTES + node.pipeMeasurementsAndValuesSerializedSize(); + size += Long.BYTES + node.serializedMeasurementsAndValuesSize(); } return size + insertRowNodeIndexList.size() * Integer.BYTES; } 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 c8a979f7e2f11..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 @@ -569,21 +569,23 @@ void subSerialize(DataOutputStream stream) throws IOException { @Override protected int serializedAttributesSize() { - return PlanNodeType.BYTES + baseSubSerializedSizeForPipe(); + return PlanNodeType.BYTES + serializedSubAttributesSize(); } /** - * Returns the exact size written by {@link #subSerialize(DataOutputStream)}. + * 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 baseSubSerializedSizeForPipe() { + final int serializedSubAttributesSize() { int size = ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()); - size += Integer.BYTES; - size += Byte.BYTES; + 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; @@ -600,16 +602,16 @@ final int baseSubSerializedSizeForPipe() { } } - size += Integer.BYTES; - size += rowCount * Long.BYTES; + size += Integer.BYTES; // row count + size += rowCount * Long.BYTES; // timestamps - size += Byte.BYTES; + 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; + size += Byte.BYTES; // whether the current measurement has a bitmap if (getBitMapIfPresent(i) != null) { size += BitMap.getSizeOfBytes(rowCount); } @@ -618,44 +620,34 @@ final int baseSubSerializedSizeForPipe() { for (int i = 0; columns != null && i < columns.length; i++) { if (shouldSerializeMeasurement(i)) { - size += columnSerializedSizeForPipe(dataTypes[i], columns[i]); + size += serializedColumnSize(dataTypes[i], columns[i]); } } - return size + Byte.BYTES; + return size + Byte.BYTES; // isAligned } - private int columnSerializedSizeForPipe(final TSDataType dataType, final Object column) { - switch (dataType) { - case INT32: - case DATE: - return rowCount * Integer.BYTES; - case INT64: - case TIMESTAMP: - return rowCount * Long.BYTES; - case FLOAT: - return rowCount * Float.BYTES; - case DOUBLE: - return rowCount * Double.BYTES; - case BOOLEAN: - return rowCount * Byte.BYTES; - case TEXT: - case BLOB: - case STRING: - case OBJECT: - int size = 0; - final Binary[] binaryValues = (Binary[]) column; - for (int i = 0; i < rowCount; i++) { - final Binary binary = binaryValues[i]; - size += - binary == null || binary.getValues() == null - ? Integer.BYTES - : Integer.BYTES + binary.getValues().length; - } - return size; - default: - throw new UnSupportedDataTypeException(String.format(DATATYPE_UNSUPPORTED, dataType)); + 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 */ 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 ae5abf83377cc..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 @@ -238,8 +238,8 @@ void subSerialize(DataOutputStream stream) throws IOException { } @Override - protected int pipeSubSerializedSize() { - return super.pipeSubSerializedSize() + getValidMeasurementNumber() * Byte.BYTES; + protected int serializedSubAttributesSize() { + return super.serializedSubAttributesSize() + getValidMeasurementNumber() * Byte.BYTES; } @Override 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 index ebaafe4fbcd2c..e26275fb8e9b2 100644 --- 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 @@ -40,6 +40,7 @@ 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; @@ -58,58 +59,101 @@ public class PipeTransferSerializationSizeTest { public void testTabletRequestLengths() throws Exception { final Tablet tablet = createTablet(); final String database = "pipe_db"; - Assert.assertEquals( - PipeTransferTabletRawReq.calculateSerializedSize(tablet), - PipeTransferTabletRawReq.toTPipeTransferReq(tablet, false).getBody().length); - Assert.assertEquals( - PipeTransferTabletRawReqV2.calculateSerializedSize(tablet, database), - PipeTransferTabletRawReqV2.toTPipeTransferReq(tablet, false, database).getBody().length); + 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 = ByteBuffer.wrap(new byte[] {1, 2, 3, 4}); - final String database = "pipe_db"; - Assert.assertEquals( - PipeTransferTabletBinaryReqV2.calculateSerializedSize(payload, database), - PipeTransferTabletBinaryReqV2.toTPipeTransferReq(payload, database).getBody().length); + 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), - PipeTransferTabletBinaryReqV2.toTPipeTransferBytes(payload, database).length); + 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), - PipeTransferTabletBinaryReq.toTPipeTransferBytes(payload).length); + 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 = ByteBuffer.wrap(new byte[] {1, 2}); - final ByteBuffer tablet = ByteBuffer.wrap(new byte[] {3, 4, 5}); - Assert.assertEquals( + 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)), - PipeTransferTabletBatchReq.toTPipeTransferReq( - Collections.singletonList(insertNode), Collections.singletonList(tablet)) - .getBody() - .length); - - final String database = "db"; - Assert.assertEquals( + 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)), - PipeTransferTabletBatchReqV2.toTPipeTransferReq( - Collections.singletonList(insertNode), - Collections.singletonList(tablet), - Collections.singletonList(database), - Collections.singletonList(database)) - .getBody() - .length); + 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 @@ -184,26 +228,27 @@ private static void assertInsertNodeRequestSizes( final ByteBuffer serializedInsertNode = insertNode.serializeToByteBuffer(); Assert.assertEquals(insertNode.serializeToByteBufferSize(), serializedInsertNode.capacity()); Assert.assertEquals(insertNode.serializeToByteBufferSize(), serializedInsertNode.remaining()); - Assert.assertEquals( - PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode), - PipeTransferTabletInsertNodeReq.toTPipeTransferReq(insertNode).getBody().length); + final PipeTransferTabletInsertNodeReq insertNodeReq = + PipeTransferTabletInsertNodeReq.toTPipeTransferReq(insertNode); + assertSerializedBodySize( + PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode), insertNodeReq.body); Assert.assertEquals( PipeTransferTabletInsertNodeReq.calculateAirGapSerializedSize(insertNode), PipeTransferTabletInsertNodeReq.toTPipeTransferBytes(insertNode).length); - Assert.assertEquals( + final PipeTransferTabletInsertNodeReqV2 insertNodeReqV2 = + PipeTransferTabletInsertNodeReqV2.toTPipeTransferReq(insertNode, databaseName); + assertSerializedBodySize( PipeTransferTabletInsertNodeReqV2.calculateSerializedSize(insertNode, databaseName), - PipeTransferTabletInsertNodeReqV2.toTPipeTransferReq(insertNode, databaseName) - .getBody() - .length); + insertNodeReqV2.body); Assert.assertEquals( PipeTransferTabletInsertNodeReqV2.calculateAirGapSerializedSize(insertNode, databaseName), PipeTransferTabletInsertNodeReqV2.toTPipeTransferBytes(insertNode, databaseName).length); - Assert.assertEquals( - IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode), + final IoTConsensusV2TabletInsertNodeReq iotConsensusReq = IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq( - insertNode, null, null, MinimumProgressIndex.INSTANCE, 0) - .getBody() - .length); + insertNode, null, null, MinimumProgressIndex.INSTANCE, 0); + assertSerializedBodySize( + IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode), + iotConsensusReq.body); } private static List indexes(final int size) { @@ -445,6 +490,33 @@ private static TsTableColumnCategory[] columnCategories() { 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( From 0a4dacd343cd283022e07b0e7d9c4ad73456f3e0 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 22 Jul 2026 11:18:55 +0800 Subject: [PATCH 4/5] Optimize Pipe tablet batch memory accounting --- .../PipeInsertNodeTabletInsertionEvent.java | 61 +++++++++++-------- .../batch/PipeTabletEventPlainBatch.java | 11 ++-- 2 files changed, 43 insertions(+), 29 deletions(-) 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/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 832287c7148ea..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 @@ -177,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( @@ -229,10 +234,6 @@ 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; } From 3fd4b7d6c20fecbe226445e96b9f4d359f1a6b53 Mon Sep 17 00:00:00 2001 From: luoluoyuyu Date: Wed, 22 Jul 2026 18:09:02 +0800 Subject: [PATCH 5/5] Optimize Pipe tablet processing memory usage --- ...sertionEventTableParserTabletIterator.java | 48 ++++++++++--------- .../memory/InsertNodeMemoryEstimator.java | 16 ++----- .../request/PipeTransferTabletBatchReqV2.java | 3 +- .../dataregion/memtable/TsFileProcessor.java | 18 +++++++ 4 files changed, 49 insertions(+), 36 deletions(-) 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/request/PipeTransferTabletBatchReqV2.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java index c1d388031bb4f..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 @@ -61,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<>(); 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 {