This is an automated email from the ASF dual-hosted git repository.

jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 8c68dd71f86 Optimize pipe request serialization buffer sizing (#18233)
8c68dd71f86 is described below

commit 8c68dd71f8685a5e408bc0eeeb9865fbe1d02cdb
Author: Zhenyu Luo <[email protected]>
AuthorDate: Tue Jul 28 10:02:58 2026 +0800

    Optimize pipe request serialization buffer sizing (#18233)
    
    * Optimize pipe request serialization buffer sizing
    
    * Add exact insert node serialization sizing
    
    * Optimize Pipe serialization buffer sizing
    
    * Optimize Pipe tablet batch memory accounting
    
    * Optimize Pipe tablet processing memory usage
    
    * fix
    
    * fix
---
 .../tablet/PipeInsertNodeTabletInsertionEvent.java |  61 ++-
 ...ileInsertionEventTableParserTabletIterator.java |  48 +-
 .../resource/memory/InsertNodeMemoryEstimator.java |  16 +-
 .../evolvable/batch/PipeTabletEventPlainBatch.java |  24 +-
 .../request/PipeTransferTabletBatchReq.java        |  27 +-
 .../request/PipeTransferTabletBatchReqV2.java      |  43 +-
 .../request/PipeTransferTabletBinaryReq.java       |  14 +-
 .../request/PipeTransferTabletBinaryReqV2.java     |  28 +-
 .../request/PipeTransferTabletInsertNodeReq.java   |  20 +-
 .../request/PipeTransferTabletInsertNodeReqV2.java |  16 +-
 .../request/PipeTransferTabletRawReq.java          |  13 +-
 .../request/PipeTransferTabletRawReqV2.java        |  14 +-
 .../IoTConsensusV2TransferBatchReqBuilder.java     |  12 +-
 .../request/IoTConsensusV2TabletInsertNodeReq.java |   6 +
 .../plan/planner/plan/PlanFragment.java            |   4 +
 .../plan/node/pipe/PipeEnrichedInsertNode.java     |  16 +
 .../plan/node/write/InsertMultiTabletsNode.java    |   9 +
 .../plan/planner/plan/node/write/InsertNode.java   |  45 ++
 .../planner/plan/node/write/InsertRowNode.java     |  67 +++
 .../planner/plan/node/write/InsertRowsNode.java    |  15 +
 .../plan/node/write/InsertRowsOfOneDeviceNode.java |  10 +
 .../planner/plan/node/write/InsertTabletNode.java  |  83 ++++
 .../plan/node/write/RelationalInsertRowNode.java   |   5 +
 .../node/write/RelationalInsertTabletNode.java     |  11 +
 .../scheduler/FragmentInstanceDispatcherImpl.java  |   2 +
 .../dataregion/memtable/TsFileProcessor.java       |  18 +
 .../request/PipeTransferSerializationSizeTest.java | 530 +++++++++++++++++++++
 .../node/write/InsertRowsNodeSerdeTest.java        |  26 +
 .../plan/planner/plan/node/PlanNode.java           |   8 +
 29 files changed, 1090 insertions(+), 101 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 5fab6979e29..3bebbc91153 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<PipeTabletMemoryBlock> allocatedMemoryBlock;
   private volatile List<Tablet> tablets;
+  // Calculated together with tablets so downstream batching does not rescan 
Tablet internals.
+  private volatile long tabletsMemoryUsageInBytes;
 
   private List<TabletInsertionEventParser> eventParsers;
 
@@ -481,22 +483,31 @@ public class PipeInsertNodeTabletInsertionEvent extends 
PipeInsertionEvent
   // TODO: for table model insertion, we need to get the database name
   public synchronized List<Tablet> convertToTablets() {
     if (Objects.isNull(tablets)) {
-      tablets =
-          initEventParsers().stream()
-              .map(TabletInsertionEventParser::convertToTablet)
-              .collect(Collectors.toList());
+      final List<TabletInsertionEventParser> parsers = initEventParsers();
+      final List<Tablet> 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<TabletInsertionEventParser> initEventParsers() {
@@ -505,35 +516,27 @@ public class PipeInsertNodeTabletInsertionEvent extends 
PipeInsertionEvent
         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 @@ public class PipeInsertNodeTabletInsertionEvent extends 
PipeInsertionEvent
     }
   }
 
+  private static int getEventParserCount(final InsertNode node) {
+    if (node instanceof InsertRowsNode) {
+      return ((InsertRowsNode) node).getInsertRowNodeList().size();
+    }
+    if (node instanceof RelationalInsertRowsNode) {
+      return ((RelationalInsertRowsNode) node).getInsertRowNodeList().size();
+    }
+    return 1;
+  }
+
   public long count() {
     long count = 0;
     for (final Tablet covertedTablet : convertToTablets()) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/parser/table/TsFileInsertionEventTableParserTabletIterator.java
index 5b50eb166be..384f475a359 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 org.apache.tsfile.write.schema.MeasurementSchema;
 
 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<Tablet> {
 
@@ -137,9 +137,12 @@ public class TsFileInsertionEventTableParserTabletIterator 
implements Iterator<T
     this.metadataQuerier = new MetadataQuerierByFileImpl(reader);
     fileMetadata = this.metadataQuerier.getWholeFileMetadata();
     final List<Map.Entry<String, TableSchema>> tableSchemaList =
-        fileMetadata.getTableSchemaMap().entrySet().stream()
-            .filter(predicate)
-            .collect(Collectors.toList());
+        new ArrayList<>(fileMetadata.getTableSchemaMap().size());
+    for (final Map.Entry<String, TableSchema> entry : 
fileMetadata.getTableSchemaMap().entrySet()) {
+      if (predicate.test(entry)) {
+        tableSchemaList.add(entry);
+      }
+    }
 
     this.allocatedMemoryBlockForTablet = allocatedMemoryBlockForTablet;
     this.allocatedMemoryBlockForBatchData = allocatedMemoryBlockForBatchData;
@@ -250,10 +253,10 @@ public class 
TsFileInsertionEventTableParserTabletIterator implements Iterator<T
               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 @@ public class 
TsFileInsertionEventTableParserTabletIterator implements Iterator<T
     timeChunk.getData().rewind();
     long size = timeChunkSize;
 
-    final List<Chunk> valueChunkList = new ArrayList<>();
+    final int fieldSchemaSize = fieldSchemaList.size();
+    final List<Chunk> valueChunkList = new ArrayList<>(fieldSchemaSize);
     final Map<String, IChunkMetadata> 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<ColumnCategory> categories = new ArrayList<>(deviceIdSize);
+    final List<ColumnCategory> 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 f4c9442ce01..394d67fd6be 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
@@ -760,24 +760,16 @@ public class InsertNodeMemoryEstimator {
     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<Integer> 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<Integer, TSStatus> results) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventPlainBatch.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/batch/PipeTabletEventPlainBatch.java
index bdf6ee1874a..3eec34cc94b 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 class PipeTabletEventPlainBatch extends 
PipeTabletEventBatch {
           }
         }
         for (final Pair<Boolean, Tablet> tabletPair : batchTablets) {
-          try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+          try (final PublicBAOS byteArrayOutputStream =
+                  new 
PublicBAOS(calculateTabletSerializedSize(tabletPair.getRight()));
               final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
             tabletPair.getRight().serialize(outputStream);
             ReadWriteIOUtils.write(true, outputStream);
@@ -176,7 +177,12 @@ public class PipeTabletEventPlainBatch extends 
PipeTabletEventBatch {
         insertNodeDataBases.add(databaseName);
       } else {
         final List<Tablet> tablets = 
pipeInsertNodeTabletInsertionEvent.convertToTablets();
-        estimateSize = calculateTabletsSizeInBytes(tablets);
+        // convertToTablets() has already measured every tablet for the event 
memory block. Reuse
+        // that exact measurement instead of calling Tablet.ramBytesUsed() 
(which walks the schema
+        // map) once more while building this batch.
+        estimateSize =
+            pipeInsertNodeTabletInsertionEvent.getTabletsMemoryUsageInBytes()
+                + (long) Integer.BYTES * tablets.size();
         increaseTotalBufferSizeAndUpdateMemoryBlock(estimateSize);
         for (final Tablet tablet : tablets) {
           constructTabletBatchWithoutMemoryReservation(
@@ -192,9 +198,11 @@ public class PipeTabletEventPlainBatch extends 
PipeTabletEventBatch {
                 pipeRawTabletInsertionEvent.convertToTablet(),
                 pipeRawTabletInsertionEvent.getTableModelDatabaseName());
       } else {
-        try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+        final Tablet tablet = pipeRawTabletInsertionEvent.convertToTablet();
+        try (final PublicBAOS byteArrayOutputStream =
+                new PublicBAOS(calculateTabletSerializedSize(tablet));
             final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
-          
pipeRawTabletInsertionEvent.convertToTablet().serialize(outputStream);
+          tablet.serialize(outputStream);
           ReadWriteIOUtils.write(pipeRawTabletInsertionEvent.isAligned(), 
outputStream);
           buffer = ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, 
byteArrayOutputStream.size());
         }
@@ -226,14 +234,14 @@ public class PipeTabletEventPlainBatch extends 
PipeTabletEventBatch {
     currentBatch.getRight().add(tablet);
   }
 
-  private long calculateTabletsSizeInBytes(final List<Tablet> tablets) {
-    return 
tablets.stream().mapToLong(PipeTabletEventPlainBatch::calculateTabletSizeInBytes).sum();
-  }
-
   private static long calculateTabletSizeInBytes(final Tablet tablet) {
     return PipeMemoryWeightUtil.calculateTabletSizeInBytes(tablet) + 4;
   }
 
+  private static int calculateTabletSerializedSize(final Tablet tablet) {
+    return tablet.serializedSize() + Byte.BYTES;
+  }
+
   static boolean mayAppendTablet(final Tablet target, final Tablet source) {
     // Tablet.append already checks schemas and column categories. Avoid 
repeating those potentially
     // expensive comparisons here because wide-table pipe transfer can have 
many columns.
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReq.java
index 352ff0bfc63..4de77179a7c 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 @@ import java.util.Objects;
 
 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<PipeTransferTabletBinaryReq> binaryReqs = new 
ArrayList<>();
   private final transient List<PipeTransferTabletInsertNodeReq> insertNodeReqs 
= new ArrayList<>();
   private final transient List<PipeTransferTabletRawReq> tabletReqs = new 
ArrayList<>();
@@ -129,19 +134,28 @@ public class PipeTransferTabletBatchReq extends 
TPipeTransferReq {
 
     batchReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion();
     batchReq.type = PipeRequestType.TRANSFER_TABLET_BATCH.getType();
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(calculateSerializedSize(insertNodeBuffers, 
tabletBuffers));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
       // Binary buffer, for rolling upgrade
       ReadWriteIOUtils.write(0, outputStream);
 
+      // Insert-node and raw-tablet serializations are self-delimiting, so 
their lengths are not
+      // written separately.
       ReadWriteIOUtils.write(insertNodeBuffers.size(), outputStream);
       for (final ByteBuffer insertNodeBuffer : insertNodeBuffers) {
-        outputStream.write(insertNodeBuffer.array(), 0, 
insertNodeBuffer.limit());
+        outputStream.write(
+            insertNodeBuffer.array(),
+            insertNodeBuffer.arrayOffset() + insertNodeBuffer.position(),
+            insertNodeBuffer.remaining());
       }
 
       ReadWriteIOUtils.write(tabletBuffers.size(), outputStream);
       for (final ByteBuffer tabletBuffer : tabletBuffers) {
-        outputStream.write(tabletBuffer.array(), 0, tabletBuffer.limit());
+        outputStream.write(
+            tabletBuffer.array(),
+            tabletBuffer.arrayOffset() + tabletBuffer.position(),
+            tabletBuffer.remaining());
       }
 
       batchReq.body =
@@ -151,6 +165,13 @@ public class PipeTransferTabletBatchReq extends 
TPipeTransferReq {
     return batchReq;
   }
 
+  static int calculateSerializedSize(
+      final List<ByteBuffer> insertNodeBuffers, final List<ByteBuffer> 
tabletBuffers) {
+    return BATCH_REQUEST_COUNT_SERIALIZED_SIZE
+        + insertNodeBuffers.stream().mapToInt(ByteBuffer::remaining).sum()
+        + tabletBuffers.stream().mapToInt(ByteBuffer::remaining).sum();
+  }
+
   public static PipeTransferTabletBatchReq fromTPipeTransferReq(
       final TPipeTransferReq transferReq) {
     final PipeTransferTabletBatchReq batchReq = new 
PipeTransferTabletBatchReq();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
index 4dd044cf273..c79cab7d88f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
@@ -46,6 +46,12 @@ import java.util.Map;
 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<PipeTransferTabletInsertNodeReqV2> 
insertNodeReqs =
       new ArrayList<>();
   private final transient List<PipeTransferTabletRawReqV2> tabletReqs = new 
ArrayList<>();
@@ -55,7 +61,8 @@ public class PipeTransferTabletBatchReqV2 extends 
TPipeTransferReq {
   }
 
   public List<InsertBaseStatement> constructStatements() {
-    final List<InsertBaseStatement> statements = new ArrayList<>();
+    final List<InsertBaseStatement> statements =
+        new ArrayList<>(insertNodeReqs.size() + tabletReqs.size());
 
     final Map<String, List<InsertRowStatement>> 
tableModelDatabaseInsertRowStatementMap =
         new LinkedHashMap<>();
@@ -186,22 +193,33 @@ public class PipeTransferTabletBatchReqV2 extends 
TPipeTransferReq {
 
     batchReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion();
     batchReq.type = PipeRequestType.TRANSFER_TABLET_BATCH_V2.getType();
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(
+                calculateSerializedSize(
+                    insertNodeBuffers, tabletBuffers, insertNodeDataBases, 
tabletDataBases));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
       // Binary buffer, for rolling upgrade
       ReadWriteIOUtils.write(0, outputStream);
 
+      // Insert-node and raw-tablet serializations are self-delimiting, so 
their lengths are not
+      // written separately.
       ReadWriteIOUtils.write(insertNodeBuffers.size(), outputStream);
       for (int i = 0; i < insertNodeBuffers.size(); i++) {
         final ByteBuffer insertNodeBuffer = insertNodeBuffers.get(i);
-        outputStream.write(insertNodeBuffer.array(), 0, 
insertNodeBuffer.limit());
+        outputStream.write(
+            insertNodeBuffer.array(),
+            insertNodeBuffer.arrayOffset() + insertNodeBuffer.position(),
+            insertNodeBuffer.remaining());
         ReadWriteIOUtils.write(insertNodeDataBases.get(i), outputStream);
       }
 
       ReadWriteIOUtils.write(tabletBuffers.size(), outputStream);
       for (int i = 0; i < tabletBuffers.size(); i++) {
         final ByteBuffer tabletBuffer = tabletBuffers.get(i);
-        outputStream.write(tabletBuffer.array(), 0, tabletBuffer.limit());
+        outputStream.write(
+            tabletBuffer.array(),
+            tabletBuffer.arrayOffset() + tabletBuffer.position(),
+            tabletBuffer.remaining());
         ReadWriteIOUtils.write(tabletDataBases.get(i), outputStream);
       }
 
@@ -212,6 +230,23 @@ public class PipeTransferTabletBatchReqV2 extends 
TPipeTransferReq {
     return batchReq;
   }
 
+  static int calculateSerializedSize(
+      final List<ByteBuffer> insertNodeBuffers,
+      final List<ByteBuffer> tabletBuffers,
+      final List<String> insertNodeDataBases,
+      final List<String> tabletDataBases) {
+    int size = BATCH_REQUEST_COUNT_SERIALIZED_SIZE;
+    for (int i = 0; i < insertNodeBuffers.size(); i++) {
+      size += insertNodeBuffers.get(i).remaining();
+      size += ReadWriteIOUtils.sizeToWrite(insertNodeDataBases.get(i));
+    }
+    for (int i = 0; i < tabletBuffers.size(); i++) {
+      size += tabletBuffers.get(i).remaining();
+      size += ReadWriteIOUtils.sizeToWrite(tabletDataBases.get(i));
+    }
+    return size;
+  }
+
   public static PipeTransferTabletBatchReqV2 fromTPipeTransferReq(
       final org.apache.iotdb.service.rpc.thrift.TPipeTransferReq transferReq) {
     final PipeTransferTabletBatchReqV2 batchReq = new 
PipeTransferTabletBatchReqV2();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReq.java
index b31816c1fcd..cf4a746cb0e 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.queryengine.plan.statement.crud.InsertBaseStatement;
 import org.apache.iotdb.db.storageengine.dataregion.wal.buffer.WALEntry;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 
-import org.apache.tsfile.utils.BytesUtils;
 import org.apache.tsfile.utils.PublicBAOS;
 import org.apache.tsfile.utils.ReadWriteIOUtils;
 
@@ -102,14 +101,23 @@ public class PipeTransferTabletBinaryReq extends 
TPipeTransferReq {
   /////////////////////////////// Air Gap ///////////////////////////////
 
   public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer) 
throws IOException {
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(calculateSerializedSize(byteBuffer));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
       ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), 
outputStream);
       ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_BINARY.getType(), 
outputStream);
-      return BytesUtils.concatByteArray(byteArrayOutputStream.toByteArray(), 
byteBuffer.array());
+      outputStream.write(
+          byteBuffer.array(),
+          byteBuffer.arrayOffset() + byteBuffer.position(),
+          byteBuffer.remaining());
+      return byteArrayOutputStream.toByteArray();
     }
   }
 
+  static int calculateSerializedSize(final ByteBuffer byteBuffer) {
+    return Byte.BYTES + Short.BYTES + byteBuffer.remaining();
+  }
+
   /////////////////////////////// Object ///////////////////////////////
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
index 2788033be2d..196af0b16d2 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBinaryReqV2.java
@@ -119,10 +119,14 @@ public class PipeTransferTabletBinaryReqV2 extends 
PipeTransferTabletBinaryReq {
 
     req.version = IoTDBSinkRequestVersion.VERSION_1.getVersion();
     req.type = PipeRequestType.TRANSFER_TABLET_BINARY_V2.getType();
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(calculateSerializedSize(byteBuffer, dataBaseName));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
-      ReadWriteIOUtils.write(byteBuffer.limit(), outputStream);
-      outputStream.write(byteBuffer.array(), 0, byteBuffer.limit());
+      ReadWriteIOUtils.write(byteBuffer.remaining(), outputStream);
+      outputStream.write(
+          byteBuffer.array(),
+          byteBuffer.arrayOffset() + byteBuffer.position(),
+          byteBuffer.remaining());
       ReadWriteIOUtils.write(dataBaseName, outputStream);
       req.body = ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, 
byteArrayOutputStream.size());
     }
@@ -151,17 +155,29 @@ public class PipeTransferTabletBinaryReqV2 extends 
PipeTransferTabletBinaryReq {
 
   public static byte[] toTPipeTransferBytes(final ByteBuffer byteBuffer, final 
String dataBaseName)
       throws IOException {
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(calculateAirGapSerializedSize(byteBuffer, 
dataBaseName));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
       ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), 
outputStream);
       
ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_BINARY_V2.getType(), 
outputStream);
-      ReadWriteIOUtils.write(byteBuffer.limit(), outputStream);
-      outputStream.write(byteBuffer.array(), 0, byteBuffer.limit());
+      ReadWriteIOUtils.write(byteBuffer.remaining(), outputStream);
+      outputStream.write(
+          byteBuffer.array(),
+          byteBuffer.arrayOffset() + byteBuffer.position(),
+          byteBuffer.remaining());
       ReadWriteIOUtils.write(dataBaseName, outputStream);
       return byteArrayOutputStream.toByteArray();
     }
   }
 
+  static int calculateSerializedSize(final ByteBuffer byteBuffer, final String 
dataBaseName) {
+    return Integer.BYTES + byteBuffer.remaining() + 
ReadWriteIOUtils.sizeToWrite(dataBaseName);
+  }
+
+  static int calculateAirGapSerializedSize(final ByteBuffer byteBuffer, final 
String dataBaseName) {
+    return Byte.BYTES + Short.BYTES + calculateSerializedSize(byteBuffer, 
dataBaseName);
+  }
+
   /////////////////////////////// Object ///////////////////////////////
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java
index bc42630d79b..42f353c2506 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletInsertNodeReq.java
@@ -31,7 +31,6 @@ import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertTablet
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
 import org.apache.iotdb.service.rpc.thrift.TPipeTransferReq;
 
-import org.apache.tsfile.utils.BytesUtils;
 import org.apache.tsfile.utils.PublicBAOS;
 import org.apache.tsfile.utils.ReadWriteIOUtils;
 
@@ -106,15 +105,28 @@ public class PipeTransferTabletInsertNodeReq extends 
TPipeTransferReq {
   /////////////////////////////// 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 b9d5eb7de85..4fc18398dda 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 class PipeTransferTabletInsertNodeReqV2 extends 
PipeTransferTabletInsertN
 
     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 class PipeTransferTabletInsertNodeReqV2 extends 
PipeTransferTabletInsertN
 
   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 class PipeTransferTabletInsertNodeReqV2 extends 
PipeTransferTabletInsertN
     }
   }
 
+  static int calculateSerializedSize(final InsertNode insertNode, final String 
dataBaseName) {
+    return PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode)
+        + ReadWriteIOUtils.sizeToWrite(dataBaseName);
+  }
+
+  static int calculateAirGapSerializedSize(final InsertNode insertNode, final 
String dataBaseName) {
+    return PipeTransferTabletInsertNodeReq.calculateAirGapSerializedSize(
+        calculateSerializedSize(insertNode, dataBaseName));
+  }
+
   /////////////////////////////// Object ///////////////////////////////
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletRawReq.java
index 01c80758152..00907ed8302 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 class PipeTransferTabletRawReq extends 
TPipeTransferReq {
 
     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 class PipeTransferTabletRawReq extends 
TPipeTransferReq {
       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 class PipeTransferTabletRawReq extends 
TPipeTransferReq {
     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 3537d9ddd00..dd6941d7f7d 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 class PipeTransferTabletRawReqV2 extends 
PipeTransferTabletRawReq {
 
     tabletReq.version = IoTDBSinkRequestVersion.VERSION_1.getVersion();
     tabletReq.type = PipeRequestType.TRANSFER_TABLET_RAW_V2.getType();
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(calculateSerializedSize(tablet, dataBaseName));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
       tablet.serialize(outputStream);
       ReadWriteIOUtils.write(isAligned, outputStream);
@@ -173,7 +174,8 @@ public class PipeTransferTabletRawReqV2 extends 
PipeTransferTabletRawReq {
 
   public static byte[] toTPipeTransferBytes(
       final Tablet tablet, final boolean isAligned, final String dataBaseName) 
throws IOException {
-    try (final PublicBAOS byteArrayOutputStream = new PublicBAOS();
+    try (final PublicBAOS byteArrayOutputStream =
+            new PublicBAOS(calculateAirGapSerializedSize(tablet, 
dataBaseName));
         final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
       ReadWriteIOUtils.write(IoTDBSinkRequestVersion.VERSION_1.getVersion(), 
outputStream);
       ReadWriteIOUtils.write(PipeRequestType.TRANSFER_TABLET_RAW_V2.getType(), 
outputStream);
@@ -184,6 +186,14 @@ public class PipeTransferTabletRawReqV2 extends 
PipeTransferTabletRawReq {
     }
   }
 
+  static int calculateSerializedSize(final Tablet tablet, final String 
dataBaseName) {
+    return tablet.serializedSize() + Byte.BYTES + 
ReadWriteIOUtils.sizeToWrite(dataBaseName);
+  }
+
+  static int calculateAirGapSerializedSize(final Tablet tablet, final String 
dataBaseName) {
+    return Byte.BYTES + Short.BYTES + calculateSerializedSize(tablet, 
dataBaseName);
+  }
+
   /////////////////////////////// Object ///////////////////////////////
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/builder/IoTConsensusV2TransferBatchReqBuilder.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/builder/IoTConsensusV2TransferBatchReqBuilder.java
index fc387084a00..8c9e0299f31 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.Logger;
 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 abstract class IoTConsensusV2TransferBatchReqBuilder 
implements AutoClose
   }
 
   protected int buildTabletInsertionBuffer(TabletInsertionEvent event) throws 
WALPipeException {
-    final ByteBuffer buffer;
     final TCommitId commitId;
 
     // event instanceof PipeInsertNodeTabletInsertionEvent)
@@ -221,17 +219,15 @@ public abstract class 
IoTConsensusV2TransferBatchReqBuilder implements AutoClose
             
pipeInsertNodeTabletInsertionEvent.getCommitterKey().getRestartTimes(),
             pipeInsertNodeTabletInsertionEvent.getRebootTimes());
 
-    // Read the bytebuffer from the wal file and transfer it directly without 
serializing or
-    // deserializing if possible
     final InsertNode insertNode = 
pipeInsertNodeTabletInsertionEvent.getInsertNode();
     // IoTConsensusV2 will transfer binary data to TIoTConsensusV2TransferReq
     final ProgressIndex progressIndex = 
pipeInsertNodeTabletInsertionEvent.getProgressIndex();
-    buffer = insertNode.serializeToByteBuffer();
-    batchReqs.add(
+    final IoTConsensusV2TabletInsertNodeReq request =
         IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq(
-            insertNode, commitId, consensusGroupId, progressIndex, 
thisDataNodeId));
+            insertNode, commitId, consensusGroupId, progressIndex, 
thisDataNodeId);
+    batchReqs.add(request);
 
-    return buffer.limit();
+    return request.body.remaining();
   }
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/request/IoTConsensusV2TabletInsertNodeReq.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/iotconsensusv2/payload/request/IoTConsensusV2TabletInsertNodeReq.java
index 5f076b68ec3..af1d0a9fc5f 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 class IoTConsensusV2TabletInsertNodeReq extends 
TIoTConsensusV2TransferRe
     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 class IoTConsensusV2TabletInsertNodeReq extends 
TIoTConsensusV2TransferRe
     return req;
   }
 
+  /** Returns the exact serialized size of an InsertNode request body. */
+  public static int calculateSerializedSize(final InsertNode insertNode) {
+    return insertNode.serializeToByteBufferSize();
+  }
+
   public static IoTConsensusV2TabletInsertNodeReq 
fromTIoTConsensusV2TransferReq(
       TIoTConsensusV2TransferReq transferReq) {
     final IoTConsensusV2TabletInsertNodeReq insertNodeReq = new 
IoTConsensusV2TabletInsertNodeReq();
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java
index a1be7a44fa7..d39281f2864 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/PlanFragment.java
@@ -251,6 +251,10 @@ public class PlanFragment {
     typeProvider = null;
   }
 
+  public void clearUselessFieldsAfterRouting() {
+    planNodeTree.clearUselessFieldsAfterRouting();
+  }
+
   public void clearTypeProvider() {
     typeProvider = null;
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/pipe/PipeEnrichedInsertNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/pipe/PipeEnrichedInsertNode.java
index f6c323a1b57..0f527ab4542 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.iotdb.db.trigger.executor.TriggerFireVisitor;
 
 import org.apache.tsfile.enums.TSDataType;
 import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.utils.ReadWriteIOUtils;
 import org.apache.tsfile.write.schema.MeasurementSchema;
 
 import java.io.DataOutputStream;
@@ -163,6 +164,11 @@ public class PipeEnrichedInsertNode extends InsertNode {
     insertNode.setDataRegionReplicaSet(dataRegionReplicaSet);
   }
 
+  @Override
+  public void clearUselessFieldsAfterRouting() {
+    insertNode.clearUselessFieldsAfterRouting();
+  }
+
   @Override
   public PartialPath getTargetPath() {
     return insertNode.getTargetPath();
@@ -290,6 +296,16 @@ public class PipeEnrichedInsertNode extends InsertNode {
     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 4d1b987b892..8d484933e8f 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 @@ public class InsertMultiTabletsNode extends InsertNode {
     }
   }
 
+  @Override
+  protected int serializedAttributesSize() {
+    int size = PlanNodeType.BYTES + Integer.BYTES;
+    for (final InsertTabletNode insertTabletNode : insertTabletNodeList) {
+      size += insertTabletNode.serializedSubAttributesSize();
+    }
+    return size + parentInsertTabletNodeIndexList.size() * Integer.BYTES;
+  }
+
   @Override
   public void markAsGeneratedByPipe() {
     isGeneratedByPipe = true;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertNode.java
index 9c8e3691883..289866c6882 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 @@ package 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write;
 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.iotdb.db.storageengine.dataregion.wal.utils.WALWriteUtils;
 import org.apache.tsfile.enums.TSDataType;
 import org.apache.tsfile.exception.NotImplementedException;
 import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.utils.PublicBAOS;
 import org.apache.tsfile.utils.ReadWriteIOUtils;
 import org.apache.tsfile.write.schema.MeasurementSchema;
 
@@ -152,6 +154,12 @@ public abstract class InsertNode extends SearchNode {
     this.dataRegionReplicaSet = dataRegionReplicaSet;
   }
 
+  @Override
+  public void clearUselessFieldsAfterRouting() {
+    super.clearUselessFieldsAfterRouting();
+    setDataRegionReplicaSet(null);
+  }
+
   public PartialPath getTargetPath() {
     return targetPath;
   }
@@ -279,6 +287,43 @@ public abstract class InsertNode extends SearchNode {
         
DataNodeQueryMessages.SERIALIZEATTRIBUTES_OF_INSERTNODE_IS_NOT_IMPLEMENTED);
   }
 
+  /**
+   * Returns the exact number of bytes written by {@link 
#serializeToByteBuffer()}.
+   *
+   * @return the serialized buffer size
+   */
+  public final int serializeToByteBufferSize() {
+    // InsertNode has no children, so PlanNode.serialize only writes the child 
count here.
+    return serializedAttributesSize() + serializedPlanNodeIdSize() + 
Integer.BYTES;
+  }
+
+  @Override
+  public ByteBuffer serializeToByteBuffer() {
+    try (final PublicBAOS byteArrayOutputStream = new 
PublicBAOS(serializeToByteBufferSize());
+        final DataOutputStream outputStream = new 
DataOutputStream(byteArrayOutputStream)) {
+      serialize(outputStream);
+      return ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, 
byteArrayOutputStream.size());
+    } catch (final IOException e) {
+      throw new SerializationRunTimeException(e);
+    }
+  }
+
+  /**
+   * Returns the exact number of bytes written by the attribute serializer.
+   *
+   * @return the serialized attribute size
+   */
+  protected abstract int serializedAttributesSize();
+
+  /**
+   * Returns the exact number of bytes written by the plan node id serializer.
+   *
+   * @return the serialized plan node id size
+   */
+  protected int serializedPlanNodeIdSize() {
+    return ReadWriteIOUtils.sizeToWrite(getPlanNodeId().getId());
+  }
+
   // region Serialization methods for WAL
 
   /** Serialized size of measurement schemas, ignoring failed time series */
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowNode.java
index 5a3628ac06c..d7a09a6b49d 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
@@ -344,6 +344,73 @@ public class InsertRowNode extends InsertNode implements 
WALEntryValue, LastCach
     subSerialize(stream);
   }
 
+  @Override
+  protected int serializedAttributesSize() {
+    return PlanNodeType.BYTES + serializedSubAttributesSize();
+  }
+
+  /**
+   * Returns the exact number of bytes written by the row serializer.
+   *
+   * @return the serialized row field size
+   */
+  protected int serializedSubAttributesSize() {
+    return Long.BYTES
+        + ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath())
+        + serializedMeasurementsAndValuesSize();
+  }
+
+  /**
+   * Returns the exact number of bytes written by the measurement and value 
serializer.
+   *
+   * @return the serialized measurement and value size
+   */
+  protected int serializedMeasurementsAndValuesSize() {
+    int size = Integer.BYTES + Byte.BYTES;
+
+    for (int i = 0; measurements != null && i < measurements.length; i++) {
+      if (!shouldSerializeMeasurement(i)) {
+        continue;
+      }
+      size +=
+          measurementSchemas == null
+              ? ReadWriteIOUtils.sizeToWrite(measurements[i])
+              : measurementSchemas[i].serializedSize();
+    }
+
+    for (int i = 0; values != null && i < values.length; i++) {
+      if (!shouldSerializeMeasurement(i)) {
+        continue;
+      }
+      size += serializedValueSize(i);
+    }
+
+    return size + Byte.BYTES + Byte.BYTES;
+  }
+
+  private int serializedValueSize(final int index) {
+    final TSDataType dataType = getDataTypeIfPresent(index);
+    if (values[index] == null) {
+      return Byte.BYTES + (dataType == null ? 0 : Byte.BYTES);
+    }
+
+    if (isNeedInferType) {
+      return Byte.BYTES + 
ReadWriteIOUtils.sizeToWrite(values[index].toString());
+    }
+
+    return Byte.BYTES
+        + switch (dataType) {
+          case BOOLEAN -> Byte.BYTES;
+          case INT32, DATE -> Integer.BYTES;
+          case INT64, TIMESTAMP -> Long.BYTES;
+          case FLOAT -> Float.BYTES;
+          case DOUBLE -> Double.BYTES;
+          case TEXT, STRING, BLOB, OBJECT -> 
ReadWriteIOUtils.sizeToWrite((Binary) values[index]);
+          case VECTOR, UNKNOWN ->
+              throw new UnSupportedDataTypeException(UNSUPPORTED_DATA_TYPE + 
dataType);
+        };
+  }
+
   void subSerialize(ByteBuffer buffer) {
     ReadWriteIOUtils.write(time, buffer);
     ReadWriteIOUtils.write(targetPath.getFullPath(), buffer);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java
index 4492bf86acd..7071f9927c0 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsNode.java
@@ -167,6 +167,12 @@ public class InsertRowsNode extends InsertNode implements 
WALEntryValue {
     results.clear();
   }
 
+  @Override
+  public void clearUselessFieldsAfterRouting() {
+    super.clearUselessFieldsAfterRouting();
+    insertRowNodeList.forEach(InsertRowNode::clearUselessFieldsAfterRouting);
+  }
+
   public TSStatus[] getFailingStatus() {
     return StatusUtils.getFailingStatus(results, insertRowNodeList.size());
   }
@@ -275,6 +281,15 @@ public class InsertRowsNode extends InsertNode implements 
WALEntryValue {
     }
   }
 
+  @Override
+  protected int serializedAttributesSize() {
+    int size = PlanNodeType.BYTES + Integer.BYTES;
+    for (InsertRowNode node : insertRowNodeList) {
+      size += node.serializedSubAttributesSize();
+    }
+    return size + insertRowNodeIndexList.size() * Integer.BYTES;
+  }
+
   @Override
   public void markAsGeneratedByPipe() {
     isGeneratedByPipe = true;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java
index ccc4ca810d8..cee1e325df2 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 @@ public class InsertRowsOfOneDeviceNode extends InsertNode {
     }
   }
 
+  @Override
+  protected int serializedAttributesSize() {
+    int size =
+        PlanNodeType.BYTES + 
ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath()) + Integer.BYTES;
+    for (InsertRowNode node : insertRowNodeList) {
+      size += Long.BYTES + node.serializedMeasurementsAndValuesSize();
+    }
+    return size + insertRowNodeIndexList.size() * Integer.BYTES;
+  }
+
   @Override
   public void markAsGeneratedByPipe() {
     isGeneratedByPipe = true;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java
index 976223bf2cf..3f39ee53412 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/InsertTabletNode.java
@@ -567,6 +567,89 @@ public class InsertTabletNode extends InsertNode 
implements WALEntryValue {
     ReadWriteIOUtils.write((byte) (isAligned ? 1 : 0), stream);
   }
 
+  @Override
+  protected int serializedAttributesSize() {
+    return PlanNodeType.BYTES + serializedSubAttributesSize();
+  }
+
+  /**
+   * Returns the exact number of bytes written by {@link 
#subSerialize(DataOutputStream)}.
+   *
+   * <p>This deliberately excludes the plan-node type, id, and children. {@link
+   * InsertMultiTabletsNode} embeds tablet nodes by calling {@code 
subSerialize}, rather than their
+   * complete plan-node serialization.
+   *
+   * @return the serialized tablet field size
+   */
+  final int serializedSubAttributesSize() {
+    int size = ReadWriteIOUtils.sizeToWrite(targetPath.getFullPath());
+
+    size += Integer.BYTES; // valid measurement count
+    size += Byte.BYTES; // whether measurement schemas are serialized
+    for (int i = 0; measurements != null && i < measurements.length; i++) {
+      if (!shouldSerializeMeasurement(i)) {
+        continue;
+      }
+      size +=
+          measurementSchemas == null
+              ? ReadWriteIOUtils.sizeToWrite(measurements[i])
+              : measurementSchemas[i].serializedSize();
+    }
+
+    for (int i = 0; dataTypes != null && i < dataTypes.length; i++) {
+      if (shouldSerializeMeasurement(i)) {
+        size += TSDataType.getSerializedSize();
+      }
+    }
+
+    size += Integer.BYTES; // row count
+    size += rowCount * Long.BYTES; // timestamps
+
+    size += Byte.BYTES; // whether bitmaps are serialized
+    if (bitMaps != null) {
+      for (int i = 0; measurements != null && i < measurements.length; i++) {
+        if (!shouldSerializeMeasurement(i)) {
+          continue;
+        }
+        size += Byte.BYTES; // whether the current measurement has a bitmap
+        if (getBitMapIfPresent(i) != null) {
+          size += BitMap.getSizeOfBytes(rowCount);
+        }
+      }
+    }
+
+    for (int i = 0; columns != null && i < columns.length; i++) {
+      if (shouldSerializeMeasurement(i)) {
+        size += serializedColumnSize(dataTypes[i], columns[i]);
+      }
+    }
+
+    return size + Byte.BYTES; // isAligned
+  }
+
+  private int serializedColumnSize(final TSDataType dataType, final Object 
column) {
+    return switch (dataType) {
+      case BOOLEAN -> rowCount * Byte.BYTES;
+      case INT32, DATE -> rowCount * Integer.BYTES;
+      case INT64, TIMESTAMP -> rowCount * Long.BYTES;
+      case FLOAT -> rowCount * Float.BYTES;
+      case DOUBLE -> rowCount * Double.BYTES;
+      case TEXT, BLOB, STRING, OBJECT -> serializedBinaryColumnSize((Binary[]) 
column);
+      case VECTOR, UNKNOWN ->
+          throw new 
UnSupportedDataTypeException(String.format(DATATYPE_UNSUPPORTED, dataType));
+    };
+  }
+
+  private int serializedBinaryColumnSize(final Binary[] binaryValues) {
+    int size = 0;
+    for (int i = 0; i < rowCount; i++) {
+      final Binary binary = binaryValues[i];
+      final byte[] values = binary == null ? null : binary.getValues();
+      size += values == null ? Integer.BYTES : Integer.BYTES + values.length;
+    }
+    return size;
+  }
+
   /** Serialize measurements or measurement schemas, ignoring failed time 
series */
   private void writeMeasurementsOrSchemas(ByteBuffer buffer) {
     ReadWriteIOUtils.write(getValidMeasurementNumber(), buffer);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertRowNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertRowNode.java
index f0d0d8de7d6..e8fac511602 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
@@ -236,6 +236,11 @@ public class RelationalInsertRowNode extends InsertRowNode 
{
     }
   }
 
+  @Override
+  protected int serializedSubAttributesSize() {
+    return super.serializedSubAttributesSize() + getValidMeasurementNumber() * 
Byte.BYTES;
+  }
+
   @Override
   protected void subSerialize(IWALByteBufferView buffer) {
     super.subSerialize(buffer);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/write/RelationalInsertTabletNode.java
index d41b078cb9b..d4c373e57d7 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 @@ public class RelationalInsertTabletNode extends 
InsertTabletNode {
     }
   }
 
+  @Override
+  protected int serializedAttributesSize() {
+    int size = super.serializedAttributesSize();
+    for (int i = 0; measurements != null && i < measurements.length; i++) {
+      if (shouldSerializeMeasurement(i)) {
+        size += Byte.BYTES;
+      }
+    }
+    return size;
+  }
+
   @Override
   public void subDeserialize(ByteBuffer buffer) {
     super.subDeserialize(buffer);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
index df719a724fe..7ae9d8d6f77 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/FragmentInstanceDispatcherImpl.java
@@ -290,6 +290,8 @@ public class FragmentInstanceDispatcherImpl implements 
IFragInstanceDispatcher {
     }
 
     try {
+      shouldDispatch.forEach(instance -> 
instance.getFragment().clearUselessFieldsAfterRouting());
+
       // 2. try the dispatch
       final List<FailedFragmentInstanceWithStatus> failedInstances =
           dispatchWriteOnce(shouldDispatch);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java
index bded5e1c627..285e29b4065 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
@@ -288,6 +288,21 @@ public class TsFileProcessor {
     }
   }
 
+  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.
    *
@@ -322,6 +337,7 @@ public class TsFileProcessor {
     // recordScheduleMemoryBlockCost
     infoForMetrics[1] += System.nanoTime() - memControlStartTime;
 
+    clearDataRegionReplicaSet(insertRowNode);
     long startTime = System.nanoTime();
     WALFlushListener walFlushListener;
     try {
@@ -427,6 +443,7 @@ public class TsFileProcessor {
     // recordScheduleMemoryBlockCost
     infoForMetrics[1] += System.nanoTime() - memControlStartTime;
 
+    clearDataRegionReplicaSet(insertRowsNode);
     long startTime = System.nanoTime();
     WALFlushListener walFlushListener;
     try {
@@ -610,6 +627,7 @@ public class TsFileProcessor {
     long[] memIncrements =
         scheduleMemoryBlock(insertTabletNode, rangeList, results, 
infoForMetrics);
 
+    clearDataRegionReplicaSet(insertTabletNode);
     long startTime = System.nanoTime();
     WALFlushListener walFlushListener;
     try {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java
new file mode 100644
index 00000000000..e26275fb8e9
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferSerializationSizeTest.java
@@ -0,0 +1,530 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.iotdb.db.pipe.sink.payload.evolvable.request;
+
+import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
+import org.apache.iotdb.commons.exception.IllegalPathException;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId;
+import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory;
+import 
org.apache.iotdb.db.pipe.sink.protocol.iotconsensusv2.payload.request.IoTConsensusV2TabletInsertNodeReq;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertMultiTabletsNode;
+import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsOfOneDeviceNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertTabletNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowsNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertTabletNode;
+
+import org.apache.tsfile.enums.ColumnCategory;
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.utils.BitMap;
+import org.apache.tsfile.utils.ReadWriteIOUtils;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+public class PipeTransferSerializationSizeTest {
+
+  @Test
+  public void testTabletRequestLengths() throws Exception {
+    final Tablet tablet = createTablet();
+    final String database = "pipe_db";
+    final PipeTransferTabletRawReq rawReq =
+        PipeTransferTabletRawReq.toTPipeTransferReq(tablet, false);
+    
assertSerializedBodySize(PipeTransferTabletRawReq.calculateSerializedSize(tablet),
 rawReq.body);
+    final PipeTransferTabletRawReqV2 rawReqV2 =
+        PipeTransferTabletRawReqV2.toTPipeTransferReq(tablet, false, database);
+    assertSerializedBodySize(
+        PipeTransferTabletRawReqV2.calculateSerializedSize(tablet, database), 
rawReqV2.body);
+    Assert.assertEquals(
+        PipeTransferTabletRawReq.calculateAirGapSerializedSize(tablet),
+        PipeTransferTabletRawReq.toTPipeTransferBytes(tablet, false).length);
+    Assert.assertEquals(
+        PipeTransferTabletRawReqV2.calculateAirGapSerializedSize(tablet, 
database),
+        PipeTransferTabletRawReqV2.toTPipeTransferBytes(tablet, false, 
database).length);
+  }
+
+  @Test
+  public void testBinaryRequestLengths() throws Exception {
+    final ByteBuffer payload = createByteBufferWithOffsetAndPosition();
+    final byte[] expectedPayload = getRemainingBytes(payload);
+    final String database = "pipe_db_\u6d4b\u8bd5";
+    final PipeTransferTabletBinaryReqV2 binaryReqV2 =
+        PipeTransferTabletBinaryReqV2.toTPipeTransferReq(payload, database);
+    assertSerializedBodySize(
+        PipeTransferTabletBinaryReqV2.calculateSerializedSize(payload, 
database), binaryReqV2.body);
+    final ByteBuffer thriftBody = binaryReqV2.body.duplicate();
+    Assert.assertEquals(expectedPayload.length, 
ReadWriteIOUtils.readInt(thriftBody));
+    assertNextBytes(thriftBody, expectedPayload);
+    Assert.assertEquals(database, ReadWriteIOUtils.readString(thriftBody));
+    Assert.assertFalse(thriftBody.hasRemaining());
+
+    final byte[] airGapV2Bytes =
+        PipeTransferTabletBinaryReqV2.toTPipeTransferBytes(payload, database);
+    Assert.assertEquals(
+        PipeTransferTabletBinaryReqV2.calculateAirGapSerializedSize(payload, 
database),
+        airGapV2Bytes.length);
+    final ByteBuffer airGapV2Body = ByteBuffer.wrap(airGapV2Bytes);
+    airGapV2Body.position(Byte.BYTES + Short.BYTES);
+    Assert.assertEquals(expectedPayload.length, 
ReadWriteIOUtils.readInt(airGapV2Body));
+    assertNextBytes(airGapV2Body, expectedPayload);
+    Assert.assertEquals(database, ReadWriteIOUtils.readString(airGapV2Body));
+    Assert.assertFalse(airGapV2Body.hasRemaining());
+
+    final byte[] airGapV1Bytes = 
PipeTransferTabletBinaryReq.toTPipeTransferBytes(payload);
+    Assert.assertEquals(
+        PipeTransferTabletBinaryReq.calculateSerializedSize(payload), 
airGapV1Bytes.length);
+    final ByteBuffer airGapV1Body = ByteBuffer.wrap(airGapV1Bytes);
+    airGapV1Body.position(Byte.BYTES + Short.BYTES);
+    assertNextBytes(airGapV1Body, expectedPayload);
+    Assert.assertFalse(airGapV1Body.hasRemaining());
+  }
+
+  @Test
+  public void testBatchRequestLengths() throws Exception {
+    final ByteBuffer insertNode = createByteBufferWithOffsetAndPosition();
+    final ByteBuffer tablet = createByteBufferWithOffsetAndPosition();
+    final byte[] expectedInsertNode = getRemainingBytes(insertNode);
+    final byte[] expectedTablet = getRemainingBytes(tablet);
+    final PipeTransferTabletBatchReq batchReq =
+        PipeTransferTabletBatchReq.toTPipeTransferReq(
+            Collections.singletonList(insertNode), 
Collections.singletonList(tablet));
+    assertSerializedBodySize(
+        PipeTransferTabletBatchReq.calculateSerializedSize(
+            Collections.singletonList(insertNode), 
Collections.singletonList(tablet)),
+        batchReq.body);
+    final ByteBuffer batchBody = batchReq.body.duplicate();
+    Assert.assertEquals(0, ReadWriteIOUtils.readInt(batchBody));
+    Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchBody));
+    assertNextBytes(batchBody, expectedInsertNode);
+    Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchBody));
+    assertNextBytes(batchBody, expectedTablet);
+    Assert.assertFalse(batchBody.hasRemaining());
+
+    final String database = "db_\u6d4b\u8bd5";
+    final PipeTransferTabletBatchReqV2 batchReqV2 =
+        PipeTransferTabletBatchReqV2.toTPipeTransferReq(
+            Collections.singletonList(insertNode),
+            Collections.singletonList(tablet),
+            Collections.singletonList(database),
+            Collections.singletonList(database));
+    assertSerializedBodySize(
+        PipeTransferTabletBatchReqV2.calculateSerializedSize(
+            Collections.singletonList(insertNode),
+            Collections.singletonList(tablet),
+            Collections.singletonList(database),
+            Collections.singletonList(database)),
+        batchReqV2.body);
+    final ByteBuffer batchV2Body = batchReqV2.body.duplicate();
+    Assert.assertEquals(0, ReadWriteIOUtils.readInt(batchV2Body));
+    Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchV2Body));
+    assertNextBytes(batchV2Body, expectedInsertNode);
+    Assert.assertEquals(database, ReadWriteIOUtils.readString(batchV2Body));
+    Assert.assertEquals(1, ReadWriteIOUtils.readInt(batchV2Body));
+    assertNextBytes(batchV2Body, expectedTablet);
+    Assert.assertEquals(database, ReadWriteIOUtils.readString(batchV2Body));
+    Assert.assertFalse(batchV2Body.hasRemaining());
+  }
+
+  @Test
+  public void testInsertNodeSerializedSize() throws Exception {
+    assertInsertNodeRequestSizes(createInsertRowNode(0), "tree_db");
+    assertInsertNodeRequestSizes(createInsertRowNodeWithSchemas(1), "tree_db");
+    assertInsertNodeRequestSizes(createInsertRowNodeWithNullValue(), 
"tree_db");
+    assertInsertNodeRequestSizes(createInsertRowNodeWithInferredType(), 
"tree_db");
+    final InsertRowNode partiallyFailedRowNode = 
createInsertRowNodeWithSchemas(2);
+    partiallyFailedRowNode.markFailedMeasurement(1);
+    assertInsertNodeRequestSizes(partiallyFailedRowNode, "tree_db");
+
+    final InsertTabletNode tabletNode = createInsertTabletNode();
+    assertInsertNodeRequestSizes(tabletNode, "tree_db");
+    assertInsertNodeRequestSizes(createInsertTabletNode(false, true), 
"tree_db");
+    final InsertTabletNode partiallyFailedTabletNode = 
createInsertTabletNode(false, true);
+    partiallyFailedTabletNode.markFailedMeasurement(1);
+    assertInsertNodeRequestSizes(partiallyFailedTabletNode, "tree_db");
+
+    final RelationalInsertRowNode relationalRowNode =
+        new RelationalInsertRowNode(
+            new PlanNodeId("relational-row"),
+            new PartialPath("table"),
+            false,
+            measurements(),
+            dataTypes(),
+            1,
+            rowValues(1),
+            false,
+            columnCategories());
+    assertInsertNodeRequestSizes(relationalRowNode, "table_db_\u6d4b\u8bd5");
+
+    final RelationalInsertTabletNode relationalTabletNode = 
createRelationalInsertTabletNode();
+    assertInsertNodeRequestSizes(relationalTabletNode, "table_db");
+
+    final List<InsertRowNode> rows = new ArrayList<>();
+    final List<InsertRowNode> relationalRows = new ArrayList<>();
+    for (int row = 0; row < 50; row++) {
+      rows.add(createInsertRowNode(row));
+      relationalRows.add(createRelationalInsertRowNode(row));
+    }
+    final InsertRowsNode insertRowsNode = new InsertRowsNode(new 
PlanNodeId("rows"));
+    insertRowsNode.setInsertRowNodeList(rows);
+    insertRowsNode.setInsertRowNodeIndexList(indexes(rows.size()));
+    assertInsertNodeRequestSizes(insertRowsNode, "tree_db");
+
+    final InsertRowsOfOneDeviceNode oneDeviceNode =
+        new InsertRowsOfOneDeviceNode(new PlanNodeId("one-device"));
+    oneDeviceNode.setInsertRowNodeList(rows);
+    oneDeviceNode.setInsertRowNodeIndexList(indexes(rows.size()));
+    assertInsertNodeRequestSizes(oneDeviceNode, "tree_db");
+
+    final InsertMultiTabletsNode multiTabletsNode =
+        new InsertMultiTabletsNode(new PlanNodeId("multi-tablets"));
+    multiTabletsNode.addInsertTabletNode(tabletNode, 0);
+    multiTabletsNode.addInsertTabletNode(relationalTabletNode, 1);
+    assertInsertNodeRequestSizes(multiTabletsNode, "tree_db");
+
+    final RelationalInsertRowsNode relationalRowsNode =
+        new RelationalInsertRowsNode(
+            new PlanNodeId("relational-rows"), indexes(relationalRows.size()), 
relationalRows);
+    assertInsertNodeRequestSizes(relationalRowsNode, "table_db");
+
+    final PipeEnrichedInsertNode pipeEnrichedInsertNode =
+        new PipeEnrichedInsertNode(createInsertRowNode(2));
+    pipeEnrichedInsertNode.setPlanNodeId(new PlanNodeId("enriched-row"));
+    assertInsertNodeRequestSizes(pipeEnrichedInsertNode, "tree_db");
+  }
+
+  private static void assertInsertNodeRequestSizes(
+      final InsertNode insertNode, final String databaseName) throws Exception 
{
+    final ByteBuffer serializedInsertNode = insertNode.serializeToByteBuffer();
+    Assert.assertEquals(insertNode.serializeToByteBufferSize(), 
serializedInsertNode.capacity());
+    Assert.assertEquals(insertNode.serializeToByteBufferSize(), 
serializedInsertNode.remaining());
+    final PipeTransferTabletInsertNodeReq insertNodeReq =
+        PipeTransferTabletInsertNodeReq.toTPipeTransferReq(insertNode);
+    assertSerializedBodySize(
+        PipeTransferTabletInsertNodeReq.calculateSerializedSize(insertNode), 
insertNodeReq.body);
+    Assert.assertEquals(
+        
PipeTransferTabletInsertNodeReq.calculateAirGapSerializedSize(insertNode),
+        
PipeTransferTabletInsertNodeReq.toTPipeTransferBytes(insertNode).length);
+    final PipeTransferTabletInsertNodeReqV2 insertNodeReqV2 =
+        PipeTransferTabletInsertNodeReqV2.toTPipeTransferReq(insertNode, 
databaseName);
+    assertSerializedBodySize(
+        PipeTransferTabletInsertNodeReqV2.calculateSerializedSize(insertNode, 
databaseName),
+        insertNodeReqV2.body);
+    Assert.assertEquals(
+        
PipeTransferTabletInsertNodeReqV2.calculateAirGapSerializedSize(insertNode, 
databaseName),
+        PipeTransferTabletInsertNodeReqV2.toTPipeTransferBytes(insertNode, 
databaseName).length);
+    final IoTConsensusV2TabletInsertNodeReq iotConsensusReq =
+        IoTConsensusV2TabletInsertNodeReq.toTIoTConsensusV2TransferReq(
+            insertNode, null, null, MinimumProgressIndex.INSTANCE, 0);
+    assertSerializedBodySize(
+        IoTConsensusV2TabletInsertNodeReq.calculateSerializedSize(insertNode),
+        iotConsensusReq.body);
+  }
+
+  private static List<Integer> indexes(final int size) {
+    final List<Integer> indexes = new ArrayList<>(size);
+    for (int i = 0; i < size; i++) {
+      indexes.add(i);
+    }
+    return indexes;
+  }
+
+  private static InsertRowNode createInsertRowNode(final int row) throws 
IllegalPathException {
+    return new InsertRowNode(
+        new PlanNodeId("row-" + row),
+        new PartialPath("root.sg.d"),
+        false,
+        measurements(),
+        dataTypes(),
+        row,
+        rowValues(row),
+        false);
+  }
+
+  private static InsertRowNode createInsertRowNodeWithSchemas(final int row)
+      throws IllegalPathException {
+    return new InsertRowNode(
+        new PlanNodeId("row-with-schemas"),
+        new PartialPath("root.sg.d"),
+        false,
+        measurements(),
+        dataTypes(),
+        measurementSchemas(),
+        row,
+        rowValues(row),
+        false);
+  }
+
+  private static InsertRowNode createInsertRowNodeWithNullValue() throws 
IllegalPathException {
+    return new InsertRowNode(
+        new PlanNodeId("row-with-null"),
+        new PartialPath("root.sg.d"),
+        false,
+        new String[] {"s"},
+        new TSDataType[] {TSDataType.INT32},
+        1,
+        new Object[] {null},
+        false);
+  }
+
+  private static InsertRowNode createInsertRowNodeWithInferredType() throws 
IllegalPathException {
+    return new InsertRowNode(
+        new PlanNodeId("row-with-inferred-type"),
+        new PartialPath("root.sg.d"),
+        false,
+        new String[] {"s"},
+        new TSDataType[] {null},
+        1,
+        new Object[] {"value"},
+        true);
+  }
+
+  private static RelationalInsertRowNode createRelationalInsertRowNode(final 
int row)
+      throws IllegalPathException {
+    return new RelationalInsertRowNode(
+        new PlanNodeId("relational-row-" + row),
+        new PartialPath("table"),
+        false,
+        measurements(),
+        dataTypes(),
+        row,
+        rowValues(row),
+        false,
+        columnCategories());
+  }
+
+  private static InsertTabletNode createInsertTabletNode() throws 
IllegalPathException {
+    return createInsertTabletNode(false, false);
+  }
+
+  private static InsertTabletNode createInsertTabletNode(final boolean 
relational)
+      throws IllegalPathException {
+    return createInsertTabletNode(relational, false);
+  }
+
+  private static InsertTabletNode createInsertTabletNode(
+      final boolean relational, final boolean withBitMaps) throws 
IllegalPathException {
+    final String[] measurements = measurements();
+    final TSDataType[] types = dataTypes();
+    final MeasurementSchema[] schemas = measurementSchemas();
+    final Object[] columns = new Object[types.length];
+    final int rowCount = 50;
+    final long[] times = new long[rowCount];
+    for (int i = 0; i < rowCount; i++) {
+      times[i] = i;
+    }
+    for (int column = 0; column < types.length; column++) {
+      switch (types[column]) {
+        case BOOLEAN:
+          final boolean[] booleanValues = new boolean[rowCount];
+          for (int row = 0; row < rowCount; row++) {
+            booleanValues[row] = row % 2 == 0;
+          }
+          columns[column] = booleanValues;
+          break;
+        case INT32:
+        case DATE:
+          final int[] intValues = new int[rowCount];
+          for (int row = 0; row < rowCount; row++) {
+            intValues[row] = row;
+          }
+          columns[column] = intValues;
+          break;
+        case INT64:
+        case TIMESTAMP:
+          final long[] longValues = new long[rowCount];
+          for (int row = 0; row < rowCount; row++) {
+            longValues[row] = row;
+          }
+          columns[column] = longValues;
+          break;
+        case FLOAT:
+          final float[] floatValues = new float[rowCount];
+          for (int row = 0; row < rowCount; row++) {
+            floatValues[row] = row;
+          }
+          columns[column] = floatValues;
+          break;
+        case DOUBLE:
+          final double[] doubleValues = new double[rowCount];
+          for (int row = 0; row < rowCount; row++) {
+            doubleValues[row] = row;
+          }
+          columns[column] = doubleValues;
+          break;
+        case TEXT:
+        case BLOB:
+        case STRING:
+        case OBJECT:
+          Binary[] values = new Binary[rowCount];
+          for (int row = 1; row < rowCount; row++) {
+            values[row] = new Binary(("value-" + 
row).getBytes(StandardCharsets.UTF_8));
+          }
+          columns[column] = values;
+          break;
+        default:
+          throw new AssertionError(types[column]);
+      }
+    }
+    final BitMap[] bitMaps = withBitMaps ? createBitMaps(types.length, 
rowCount) : null;
+    return relational
+        ? new RelationalInsertTabletNode(
+            new PlanNodeId("relational-tablet"),
+            new PartialPath("table"),
+            false,
+            measurements,
+            types,
+            schemas,
+            times,
+            bitMaps,
+            columns,
+            rowCount,
+            columnCategories())
+        : new InsertTabletNode(
+            new PlanNodeId("tablet"),
+            new PartialPath("root.sg.d"),
+            false,
+            measurements,
+            types,
+            schemas,
+            times,
+            bitMaps,
+            columns,
+            rowCount);
+  }
+
+  private static RelationalInsertTabletNode createRelationalInsertTabletNode()
+      throws IllegalPathException {
+    return (RelationalInsertTabletNode) createInsertTabletNode(true);
+  }
+
+  private static String[] measurements() {
+    return new String[] {"b", "i", "l", "f", "d", "t", "ts", "date", "blob", 
"string", "object"};
+  }
+
+  private static TSDataType[] dataTypes() {
+    return new TSDataType[] {
+      TSDataType.BOOLEAN,
+      TSDataType.INT32,
+      TSDataType.INT64,
+      TSDataType.FLOAT,
+      TSDataType.DOUBLE,
+      TSDataType.TEXT,
+      TSDataType.TIMESTAMP,
+      TSDataType.DATE,
+      TSDataType.BLOB,
+      TSDataType.STRING,
+      TSDataType.OBJECT
+    };
+  }
+
+  private static MeasurementSchema[] measurementSchemas() {
+    final String[] measurements = measurements();
+    final TSDataType[] types = dataTypes();
+    final MeasurementSchema[] schemas = new MeasurementSchema[types.length];
+    for (int i = 0; i < types.length; i++) {
+      schemas[i] = new MeasurementSchema(measurements[i], types[i], 
TSEncoding.PLAIN);
+    }
+    return schemas;
+  }
+
+  private static BitMap[] createBitMaps(final int columnCount, final int 
rowCount) {
+    final BitMap[] bitMaps = new BitMap[columnCount];
+    bitMaps[0] = new BitMap(rowCount);
+    bitMaps[0].mark(0);
+    bitMaps[columnCount - 1] = new BitMap(rowCount);
+    bitMaps[columnCount - 1].mark(rowCount - 1);
+    return bitMaps;
+  }
+
+  private static Object[] rowValues(final int row) {
+    return new Object[] {
+      true,
+      row,
+      (long) row,
+      (float) row,
+      (double) row,
+      new Binary(("text-" + row).getBytes(StandardCharsets.UTF_8)),
+      (long) row,
+      row,
+      new Binary(("blob-" + row).getBytes(StandardCharsets.UTF_8)),
+      new Binary(("string-" + row).getBytes(StandardCharsets.UTF_8)),
+      new Binary(("object-" + row).getBytes(StandardCharsets.UTF_8))
+    };
+  }
+
+  private static TsTableColumnCategory[] columnCategories() {
+    final TsTableColumnCategory[] categories = new 
TsTableColumnCategory[dataTypes().length];
+    Arrays.fill(categories, TsTableColumnCategory.FIELD);
+    categories[0] = TsTableColumnCategory.TAG;
+    return categories;
+  }
+
+  private static ByteBuffer createByteBufferWithOffsetAndPosition() {
+    final ByteBuffer source = ByteBuffer.wrap(new byte[] {0, 1, 2, 3, 4, 5});
+    source.position(1);
+    final ByteBuffer buffer = source.slice();
+    buffer.position(1);
+    buffer.limit(4);
+    return buffer;
+  }
+
+  private static byte[] getRemainingBytes(final ByteBuffer buffer) {
+    final ByteBuffer duplicate = buffer.duplicate();
+    final byte[] bytes = new byte[duplicate.remaining()];
+    duplicate.get(bytes);
+    return bytes;
+  }
+
+  private static void assertSerializedBodySize(final int expectedSize, final 
ByteBuffer body) {
+    Assert.assertEquals(expectedSize, body.remaining());
+    Assert.assertEquals(expectedSize, body.capacity());
+  }
+
+  private static void assertNextBytes(final ByteBuffer buffer, final byte[] 
expectedBytes) {
+    final byte[] actualBytes = new byte[expectedBytes.length];
+    buffer.get(actualBytes);
+    Assert.assertArrayEquals(expectedBytes, actualBytes);
+  }
+
+  private static Tablet createTablet() {
+    final Tablet tablet =
+        new Tablet(
+            "table1", Collections.singletonList(new MeasurementSchema("s1", 
TSDataType.INT32)), 1);
+    
tablet.setColumnCategories(Collections.singletonList(ColumnCategory.FIELD));
+    tablet.addTimestamp(0, 1L);
+    tablet.addValue(0, 0, 1);
+    tablet.setRowSize(1);
+    return tablet;
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java
index 907af18ef9a..e09a71e4351 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/write/InsertRowsNodeSerdeTest.java
@@ -19,11 +19,13 @@
 
 package org.apache.iotdb.db.queryengine.plan.planner.node.write;
 
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
 import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId;
 import 
org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeType;
 import org.apache.iotdb.commons.schema.table.column.TsTableColumnCategory;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowNode;
@@ -45,6 +47,30 @@ import java.nio.charset.StandardCharsets;
 
 public class InsertRowsNodeSerdeTest {
 
+  @Test
+  public void testClearUselessFieldAfterDispatch() throws IllegalPathException 
{
+    final InsertRowsNode insertRowsNode = new InsertRowsNode(new 
PlanNodeId("insert rows"));
+    final InsertRowNode insertRowNode =
+        new InsertRowNode(
+            new PlanNodeId("insert row"),
+            new PartialPath("root.sg.d1"),
+            false,
+            new String[] {"s1"},
+            new TSDataType[] {TSDataType.INT32},
+            1L,
+            new Object[] {1},
+            false);
+    insertRowsNode.addOneInsertRowNode(insertRowNode, 0);
+
+    insertRowsNode.setDataRegionReplicaSet(new TRegionReplicaSet());
+    insertRowNode.setDataRegionReplicaSet(new TRegionReplicaSet());
+
+    new 
PipeEnrichedInsertNode(insertRowsNode).clearUselessFieldsAfterRouting();
+
+    Assert.assertNull(insertRowsNode.getDataRegionReplicaSet());
+    Assert.assertNull(insertRowNode.getDataRegionReplicaSet());
+  }
+
   @Test
   public void TestSerializeAndDeserialize() throws IllegalPathException {
     InsertRowsNode node = new InsertRowsNode(new PlanNodeId("plan node 1"));
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java
index b2e74b32af8..3411efde47e 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/queryengine/plan/planner/plan/node/PlanNode.java
@@ -78,6 +78,14 @@ public abstract class PlanNode implements IConsensusRequest {
 
   public abstract void addChild(PlanNode child);
 
+  /** Releases fields that are no longer needed after the target region has 
been determined. */
+  public void clearUselessFieldsAfterRouting() {
+    final List<PlanNode> children = getChildren();
+    if (children != null) {
+      children.forEach(PlanNode::clearUselessFieldsAfterRouting);
+    }
+  }
+
   /**
    * If this plan node has to be serialized or deserialized, override this 
method. If this method is
    * overridden, the serialization and deserialization methods must be 
implemented.


Reply via email to