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.