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

JackieTien97 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 957204481d9 Fragment large TsBlocks in MPP DataExchange responses 
(#18645)
957204481d9 is described below

commit 957204481d95d83addccd7e948695c844056d816
Author: Weihao Li <[email protected]>
AuthorDate: Sat Sep 19 07:53:48 2026 +0800

    Fragment large TsBlocks in MPP DataExchange responses (#18645)
---
 .../it/env/cluster/config/MppDataNodeConfig.java   |   9 +
 .../it/env/remote/config/RemoteDataNodeConfig.java |   6 +
 .../apache/iotdb/itbase/env/DataNodeConfig.java    |   2 +
 ...AbstractIoTDBMPPDataExchangeLargeTsBlockIT.java |  98 +++++++++++
 .../IoTDBMPPDataExchangeLargeTsBlock128KIT.java    |  51 ++++++
 .../apache/iotdb/db/i18n/DataNodeMiscMessages.java |   9 +
 .../iotdb/db/i18n/DataNodeQueryMessages.java       |  20 +++
 .../apache/iotdb/db/i18n/DataNodeMiscMessages.java |   9 +
 .../iotdb/db/i18n/DataNodeQueryMessages.java       |  20 +++
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |  10 ++
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |  49 ++++++
 .../execution/exchange/MPPDataExchangeManager.java |  44 +++++
 .../execution/exchange/sink/SinkChannel.java       |  49 ++++--
 .../execution/exchange/source/SourceHandle.java    | 158 +++++++++++++++++-
 .../org/apache/iotdb/db/conf/PropertiesTest.java   |  28 ++++
 .../execution/exchange/SinkChannelTest.java        |  16 +-
 .../execution/exchange/SourceHandleTest.java       | 185 ++++++++++++++++++++-
 .../db/queryengine/execution/exchange/Utils.java   |   7 +
 .../conf/iotdb-system.properties.template          |   7 +
 .../commons/audit/UserDataTransferErrorCode.java   |   3 +-
 .../src/main/thrift/datanode.thrift                |   6 +
 21 files changed, 768 insertions(+), 18 deletions(-)

diff --git 
a/integration-test/src/main/java/org/apache/iotdb/it/env/cluster/config/MppDataNodeConfig.java
 
b/integration-test/src/main/java/org/apache/iotdb/it/env/cluster/config/MppDataNodeConfig.java
index 8c5f56ad9f1..f7dd625704c 100644
--- 
a/integration-test/src/main/java/org/apache/iotdb/it/env/cluster/config/MppDataNodeConfig.java
+++ 
b/integration-test/src/main/java/org/apache/iotdb/it/env/cluster/config/MppDataNodeConfig.java
@@ -186,4 +186,13 @@ public class MppDataNodeConfig extends MppBaseConfig 
implements DataNodeConfig {
     setProperty("table_query_device_entry_batch_size_in_bytes", 
String.valueOf(batchSizeInBytes));
     return this;
   }
+
+  @Override
+  public DataNodeConfig setMppDataExchangeMaxPayloadSizeInBytes(
+      int mppDataExchangeMaxPayloadSizeInBytes) {
+    setProperty(
+        "mpp_data_exchange_max_payload_size_in_bytes",
+        String.valueOf(mppDataExchangeMaxPayloadSizeInBytes));
+    return this;
+  }
 }
diff --git 
a/integration-test/src/main/java/org/apache/iotdb/it/env/remote/config/RemoteDataNodeConfig.java
 
b/integration-test/src/main/java/org/apache/iotdb/it/env/remote/config/RemoteDataNodeConfig.java
index c97e2b5065a..13eeb343919 100644
--- 
a/integration-test/src/main/java/org/apache/iotdb/it/env/remote/config/RemoteDataNodeConfig.java
+++ 
b/integration-test/src/main/java/org/apache/iotdb/it/env/remote/config/RemoteDataNodeConfig.java
@@ -130,4 +130,10 @@ public class RemoteDataNodeConfig implements 
DataNodeConfig {
   public DataNodeConfig setTableQueryDeviceEntryBatchSizeInBytes(long 
batchSizeInBytes) {
     return this;
   }
+
+  @Override
+  public DataNodeConfig setMppDataExchangeMaxPayloadSizeInBytes(
+      int mppDataExchangeMaxPayloadSizeInBytes) {
+    return this;
+  }
 }
diff --git 
a/integration-test/src/main/java/org/apache/iotdb/itbase/env/DataNodeConfig.java
 
b/integration-test/src/main/java/org/apache/iotdb/itbase/env/DataNodeConfig.java
index fb969c778ea..d321053f278 100644
--- 
a/integration-test/src/main/java/org/apache/iotdb/itbase/env/DataNodeConfig.java
+++ 
b/integration-test/src/main/java/org/apache/iotdb/itbase/env/DataNodeConfig.java
@@ -67,4 +67,6 @@ public interface DataNodeConfig {
   DataNodeConfig setDnMultiDirStrategy(String multiDirStrategy);
 
   DataNodeConfig setTableQueryDeviceEntryBatchSizeInBytes(long 
batchSizeInBytes);
+
+  DataNodeConfig setMppDataExchangeMaxPayloadSizeInBytes(int 
mppDataExchangeMaxPayloadSizeInBytes);
 }
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/AbstractIoTDBMPPDataExchangeLargeTsBlockIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/AbstractIoTDBMPPDataExchangeLargeTsBlockIT.java
new file mode 100644
index 00000000000..34e7798d73a
--- /dev/null
+++ 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/AbstractIoTDBMPPDataExchangeLargeTsBlockIT.java
@@ -0,0 +1,98 @@
+/*
+ * 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.relational.it.query.recent;
+
+import org.apache.iotdb.isession.ITableSession;
+import org.apache.iotdb.it.env.EnvFactory;
+
+import org.apache.tsfile.enums.ColumnCategory;
+import org.apache.tsfile.enums.TSDataType;
+import org.apache.tsfile.utils.Binary;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.IMeasurementSchema;
+import org.apache.tsfile.write.schema.MeasurementSchema;
+import org.junit.Test;
+
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.Statement;
+import java.util.List;
+
+import static org.junit.Assert.assertArrayEquals;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+public abstract class AbstractIoTDBMPPDataExchangeLargeTsBlockIT {
+
+  protected static final String DATABASE_NAME = "large_tsblock";
+  protected static final int PAYLOAD_SIZE_IN_BYTES = 128 * 1024;
+  private static final int BLOB_SIZE_IN_BYTES = 2 * PAYLOAD_SIZE_IN_BYTES + 1;
+  private static final byte[] EXPECTED_BLOB = createBlob();
+
+  protected static void prepareData() throws Exception {
+    try (Connection connection = EnvFactory.getEnv().getTableConnection();
+        Statement statement = connection.createStatement()) {
+      statement.execute("CREATE DATABASE " + DATABASE_NAME);
+      statement.execute("USE " + DATABASE_NAME);
+      statement.execute("CREATE TABLE large_blob(payload BLOB FIELD)");
+    }
+
+    List<IMeasurementSchema> schemas = List.of(new 
MeasurementSchema("payload", TSDataType.BLOB));
+    Tablet tablet =
+        new Tablet(
+            "large_blob",
+            IMeasurementSchema.getMeasurementNameList(schemas),
+            IMeasurementSchema.getDataTypeList(schemas),
+            List.of(ColumnCategory.FIELD),
+            1);
+    tablet.addTimestamp(0, 1);
+    tablet.addValue("payload", 0, new Binary(EXPECTED_BLOB));
+
+    try (ITableSession session =
+        EnvFactory.getEnv().getTableSessionConnectionWithDB(DATABASE_NAME)) {
+      session.insert(tablet);
+    }
+  }
+
+  @Test
+  public void testLargeBlobTransferredInFragments() throws Exception {
+    assertTrue(EXPECTED_BLOB.length > PAYLOAD_SIZE_IN_BYTES);
+
+    try (Connection connection = EnvFactory.getEnv().getTableConnection();
+        Statement statement = connection.createStatement()) {
+      statement.execute("USE " + DATABASE_NAME);
+      try (ResultSet resultSet = statement.executeQuery("SELECT time, payload 
FROM large_blob")) {
+        assertTrue(resultSet.next());
+        assertEquals(1, resultSet.getLong("time"));
+        assertArrayEquals(EXPECTED_BLOB, resultSet.getBytes("payload"));
+        assertFalse(resultSet.next());
+      }
+    }
+  }
+
+  private static byte[] createBlob() {
+    byte[] blob = new byte[BLOB_SIZE_IN_BYTES];
+    for (int i = 0; i < blob.length; i++) {
+      blob[i] = (byte) (i * 31 + 7);
+    }
+    return blob;
+  }
+}
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/IoTDBMPPDataExchangeLargeTsBlock128KIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/IoTDBMPPDataExchangeLargeTsBlock128KIT.java
new file mode 100644
index 00000000000..672859429f3
--- /dev/null
+++ 
b/integration-test/src/test/java/org/apache/iotdb/relational/it/query/recent/IoTDBMPPDataExchangeLargeTsBlock128KIT.java
@@ -0,0 +1,51 @@
+/*
+ * 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.relational.it.query.recent;
+
+import org.apache.iotdb.it.env.EnvFactory;
+import org.apache.iotdb.it.framework.IoTDBTestRunner;
+import org.apache.iotdb.itbase.category.TableClusterIT;
+
+import org.junit.AfterClass;
+import org.junit.BeforeClass;
+import org.junit.experimental.categories.Category;
+import org.junit.runner.RunWith;
+
+@RunWith(IoTDBTestRunner.class)
+@Category({TableClusterIT.class})
+public class IoTDBMPPDataExchangeLargeTsBlock128KIT
+    extends AbstractIoTDBMPPDataExchangeLargeTsBlockIT {
+
+  @BeforeClass
+  public static void setUp() throws Exception {
+    EnvFactory.getEnv()
+        .getConfig()
+        .getDataNodeConfig()
+        .setMppDataExchangeMaxPayloadSizeInBytes(PAYLOAD_SIZE_IN_BYTES);
+    
EnvFactory.getEnv().getConfig().getCommonConfig().setDataReplicationFactor(1);
+    EnvFactory.getEnv().initClusterEnvironment();
+    prepareData();
+  }
+
+  @AfterClass
+  public static void tearDown() {
+    EnvFactory.getEnv().cleanClusterEnvironment();
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
index b571df825ab..c6a94f23b17 100644
--- 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
@@ -1491,5 +1491,14 @@ public final class DataNodeMiscMessages {
   public static final String
       
LOG_TABLE_QUERY_DEVICE_ENTRY_BATCH_SIZE_IN_BYTES_ARG_EXCEEDS_DN_THRIFT_MAX_FRAME_SIZE_ARG_USING_ARG_AS_THE_EFFECTIVE_VALUE_2AE1BEDA
 =
           "table_query_device_entry_batch_size_in_bytes (%d) exceeds the 
maximum RPC payload (dn_thrift_max_frame_size %d minus 1024 bytes); using %d as 
the effective value";
+  public static final String
+      
LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_IS_NOT_POSITIVE_USING_DEFAULT_VALUE_ARG_1AA821B2
 =
+          "mpp_data_exchange_max_payload_size_in_bytes (%d) is not positive, 
using default value %d";
+  public static final String
+      
LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_EXCEEDS_MAXIMUM_ALLOWED_VALUE_ARG_USING_ARG_D9BF0BBC
 =
+          "mpp_data_exchange_max_payload_size_in_bytes (%d) exceeds the 
maximum allowed value %d, using %d";
+  public static final String
+      
LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_IS_BELOW_MINIMUM_ALLOWED_VALUE_ARG_USING_ARG_794ABC76
 =
+          "mpp_data_exchange_max_payload_size_in_bytes (%d) is below the 
minimum allowed value %d, using %d";
 
 }
diff --git 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
index 771e3a4517d..ef13dccd8ac 100644
--- 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
@@ -1423,6 +1423,26 @@ public final class DataNodeQueryMessages {
       "failed to get data block [{}, {}), attempt times: {}";
   public static final String 
EXCEPTION_UNEXPECTED_DATA_BLOCK_RESPONSE_SIZE_A7DD7E33 =
       "Unexpected data block response size.";
+  public static final String
+      
EXCEPTION_INVALID_SERIALIZED_TSBLOCK_FRAGMENT_OFFSET_ARG_FOR_BLOCK_SIZE_ARG_53BC0284
 =
+          "Invalid serialized TsBlock fragment offset %s for block size %s.";
+  public static final String
+      
EXCEPTION_ACCUMULATED_TSBLOCK_FRAGMENT_LENGTH_ARG_DOES_NOT_MATCH_TOTALLENGTH_ARG_1B784303
 =
+          "Accumulated TsBlock fragment length %s does not match totalLength 
%s.";
+  public static final String
+      
EXCEPTION_THE_FIRST_FRAGMENTED_DATA_BLOCK_RESPONSE_MUST_INCLUDE_TOTALLENGTH_C5C79BC2
 =
+          "The first fragmented data block response must include totalLength.";
+  public static final String
+      
EXCEPTION_NEXT_SEQUENCE_ID_ARG_EXCEEDS_REQUESTED_END_SEQUENCE_ID_ARG_30B1726E =
+          "Next sequence ID %s exceeds requested end sequence ID %s.";
+  public static final String
+      
EXCEPTION_A_COMPLETED_DATA_BLOCK_RESPONSE_RANGE_MUST_NOT_RETAIN_A_PARTIAL_TSBLOCK_85E5C287
 =
+          "A completed data block response range must not retain a partial 
TsBlock.";
+  public static final String 
EXCEPTION_TSBLOCK_FRAGMENT_MUST_NOT_BE_EMPTY_C7D19863 =
+      "TsBlock fragment must not be empty.";
+  public static final String
+      
EXCEPTION_NEXT_FRAGMENT_OFFSET_ARG_MUST_BE_GREATER_THAN_CURRENT_OFFSET_ARG_AND_MATCH_ACCUMULATED_FRAGMENT_LENGTH_ARG_ECC31047
 =
+          "Next fragment offset %s must be greater than current offset %s and 
match accumulated fragment length %s.";
   public static final String FAILED_TO_SEND_ACK_DATA_BLOCK_EVENT =
       "failed to send ack data block event [{}, {}), attempt times: {}";
   public static final String SEND_CLOSE_SINK_CHANNEL_EVENT_FAILED =
diff --git 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
index bf15d6cde07..8bc63325f8d 100644
--- 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
@@ -1471,5 +1471,14 @@ public final class DataNodeMiscMessages {
   public static final String
       
LOG_TABLE_QUERY_DEVICE_ENTRY_BATCH_SIZE_IN_BYTES_ARG_EXCEEDS_DN_THRIFT_MAX_FRAME_SIZE_ARG_USING_ARG_AS_THE_EFFECTIVE_VALUE_2AE1BEDA
 =
           "table_query_device_entry_batch_size_in_bytes(%d)超过最大 RPC 
payload(dn_thrift_max_frame_size %d 减去 1024 字节),将使用 %d 作为生效值";
+  public static final String
+      
LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_IS_NOT_POSITIVE_USING_DEFAULT_VALUE_ARG_1AA821B2
 =
+          "mpp_data_exchange_max_payload_size_in_bytes(%d)不是正数,将使用默认值 %d";
+  public static final String
+      
LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_EXCEEDS_MAXIMUM_ALLOWED_VALUE_ARG_USING_ARG_D9BF0BBC
 =
+          "mpp_data_exchange_max_payload_size_in_bytes(%d)超过允许的最大值 %d,将使用 %d";
+  public static final String
+      
LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_IS_BELOW_MINIMUM_ALLOWED_VALUE_ARG_USING_ARG_794ABC76
 =
+          "mpp_data_exchange_max_payload_size_in_bytes(%d)低于允许的最小值 %d,将使用 %d";
 
 }
diff --git 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
index 24170c487e5..f9dde6564ce 100644
--- 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
@@ -1403,6 +1403,26 @@ public final class DataNodeQueryMessages {
       "获取数据块 [{}, {}) 失败,尝试次数:{}";
   public static final String 
EXCEPTION_UNEXPECTED_DATA_BLOCK_RESPONSE_SIZE_A7DD7E33 =
       "数据块响应数量异常。";
+  public static final String
+      
EXCEPTION_INVALID_SERIALIZED_TSBLOCK_FRAGMENT_OFFSET_ARG_FOR_BLOCK_SIZE_ARG_53BC0284
 =
+          "序列化 TsBlock 分片偏移量 %s 无效,块大小为 %s。";
+  public static final String
+      
EXCEPTION_ACCUMULATED_TSBLOCK_FRAGMENT_LENGTH_ARG_DOES_NOT_MATCH_TOTALLENGTH_ARG_1B784303
 =
+          "TsBlock 分片累计长度 %s 与 totalLength %s 不一致。";
+  public static final String
+      
EXCEPTION_THE_FIRST_FRAGMENTED_DATA_BLOCK_RESPONSE_MUST_INCLUDE_TOTALLENGTH_C5C79BC2
 =
+          "首个分片数据块响应必须包含 totalLength。";
+  public static final String
+      
EXCEPTION_NEXT_SEQUENCE_ID_ARG_EXCEEDS_REQUESTED_END_SEQUENCE_ID_ARG_30B1726E =
+          "下一个 sequence ID %s 超过请求的结束 sequence ID %s。";
+  public static final String
+      
EXCEPTION_A_COMPLETED_DATA_BLOCK_RESPONSE_RANGE_MUST_NOT_RETAIN_A_PARTIAL_TSBLOCK_85E5C287
 =
+          "已完成的数据块响应区间不能保留未完成的 TsBlock。";
+  public static final String 
EXCEPTION_TSBLOCK_FRAGMENT_MUST_NOT_BE_EMPTY_C7D19863 =
+      "TsBlock 分片不能为空。";
+  public static final String
+      
EXCEPTION_NEXT_FRAGMENT_OFFSET_ARG_MUST_BE_GREATER_THAN_CURRENT_OFFSET_ARG_AND_MATCH_ACCUMULATED_FRAGMENT_LENGTH_ARG_ECC31047
 =
+          "下一分片偏移量 %s 必须大于当前偏移量 %s,并且等于分片累计长度 %s。";
   public static final String FAILED_TO_SEND_ACK_DATA_BLOCK_EVENT =
       "发送数据块确认事件 [{}, {}) 失败,尝试次数:{}";
   public static final String SEND_CLOSE_SINK_CHANNEL_EVENT_FAILED =
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index 44a9daab3a8..b3f34ea898f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -950,6 +950,8 @@ public class IoTDBConfig {
   /** Core pool size of mpp data exchange. */
   private int mppDataExchangeCorePoolSize = 10;
 
+  private int mppDataExchangeMaxPayloadSizeInBytes = 4 * 1024 * 1024;
+
   /** Max pool size of mpp data exchange. */
   private int mppDataExchangeMaxPoolSize = 10;
 
@@ -3417,6 +3419,14 @@ public class IoTDBConfig {
     this.mppDataExchangeKeepAliveTimeInMs = mppDataExchangeKeepAliveTimeInMs;
   }
 
+  public int getMppDataExchangeMaxPayloadSizeInBytes() {
+    return mppDataExchangeMaxPayloadSizeInBytes;
+  }
+
+  public void setMppDataExchangeMaxPayloadSizeInBytes(int 
mppDataExchangeMaxPayloadSizeInBytes) {
+    this.mppDataExchangeMaxPayloadSizeInBytes = 
mppDataExchangeMaxPayloadSizeInBytes;
+  }
+
   public int getConnectionTimeoutInMS() {
     return connectionTimeoutInMS;
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index 1eff0fcb123..12f4d7e906e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -119,6 +119,10 @@ public class IoTDBDescriptor {
 
   private static final double MIN_DIR_USE_PROPORTION = 0.5;
 
+  private static final int DEFAULT_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES 
= 4 * 1024 * 1024;
+
+  private static final int MIN_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES = 
128 * 1024;
+
   private static final String[] DEFAULT_WAL_THRESHOLD_NAME = {
     "iot_consensus_throttle_threshold_in_byte", 
"wal_throttle_threshold_in_byte"
   };
@@ -2277,6 +2281,8 @@ public class IoTDBDescriptor {
       memoryConfig.loadTableQueryDeviceEntryBatchSize(
           properties, conf.getThriftMaxFrameSize(), LOGGER);
 
+      loadMppDataExchangeMaxPayloadSize(properties);
+
       // update wal config
       long prevDeleteWalFilesPeriodInMs = conf.getDeleteWalFilesPeriodInMs();
       loadWALHotModifiedProps(properties);
@@ -2454,6 +2460,9 @@ public class IoTDBDescriptor {
     ConfigurationFileUtils.updateAppliedProperties(
         "table_query_device_entry_batch_size_in_bytes",
         
Long.toString(memoryConfig.getTableQueryDeviceEntryBatchSizeInBytes()));
+    ConfigurationFileUtils.updateAppliedProperties(
+        "mpp_data_exchange_max_payload_size_in_bytes",
+        Integer.toString(conf.getMppDataExchangeMaxPayloadSizeInBytes()));
     ConfigurationFileUtils.updateAppliedProperties(
         DEFAULT_WAL_THRESHOLD_NAME[1], 
Long.toString(conf.getThrottleThreshold()));
   }
@@ -3076,6 +3085,8 @@ public class IoTDBDescriptor {
                 "mpp_data_exchange_keep_alive_time_in_ms",
                 
Integer.toString(conf.getMppDataExchangeKeepAliveTimeInMs()))));
 
+    loadMppDataExchangeMaxPayloadSize(properties);
+
     conf.setPartitionCacheSize(
         Integer.parseInt(
             properties.getProperty(
@@ -3088,6 +3099,44 @@ public class IoTDBDescriptor {
                 
Integer.toString(commonConfig.getDriverTaskExecutionTimeSliceInMs()))));
   }
 
+  private void loadMppDataExchangeMaxPayloadSize(TrimProperties properties) {
+    int configuredSize =
+        Integer.parseInt(
+            properties.getProperty(
+                "mpp_data_exchange_max_payload_size_in_bytes",
+                
Integer.toString(conf.getMppDataExchangeMaxPayloadSizeInBytes())));
+    if (configuredSize <= 0) {
+      LOGGER.warn(
+          String.format(
+              DataNodeMiscMessages
+                  
.LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_IS_NOT_POSITIVE_USING_DEFAULT_VALUE_ARG_1AA821B2,
+              configuredSize,
+              DEFAULT_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES));
+      configuredSize = DEFAULT_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES;
+    } else if (configuredSize < 
MIN_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES) {
+      LOGGER.warn(
+          String.format(
+              DataNodeMiscMessages
+                  
.LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_IS_BELOW_MINIMUM_ALLOWED_VALUE_ARG_USING_ARG_794ABC76,
+              configuredSize,
+              MIN_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES,
+              MIN_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES));
+      configuredSize = MIN_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_IN_BYTES;
+    }
+    int maxAllowedSize = conf.getThriftMaxFrameSize() - 1024;
+    if (configuredSize > maxAllowedSize) {
+      LOGGER.warn(
+          String.format(
+              DataNodeMiscMessages
+                  
.LOG_MPP_DATA_EXCHANGE_MAX_PAYLOAD_SIZE_ARG_EXCEEDS_MAXIMUM_ALLOWED_VALUE_ARG_USING_ARG_D9BF0BBC,
+              configuredSize,
+              maxAllowedSize,
+              maxAllowedSize));
+      configuredSize = maxAllowedSize;
+    }
+    conf.setMppDataExchangeMaxPayloadSizeInBytes(configuredSize);
+  }
+
   /** Get default encode algorithm by data type */
   public TSEncoding getDefaultEncodingByType(TSDataType dataType) {
     switch (dataType) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/MPPDataExchangeManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/MPPDataExchangeManager.java
index ceba3880122..61e776d7d9e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/MPPDataExchangeManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/MPPDataExchangeManager.java
@@ -24,6 +24,7 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.client.IClientManager;
 import 
org.apache.iotdb.commons.client.sync.SyncDataNodeMPPDataExchangeServiceClient;
 import org.apache.iotdb.commons.utils.TestOnly;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
 import 
org.apache.iotdb.db.queryengine.exception.exchange.GetTsBlockFromClosedOrAbortedChannelException;
 import org.apache.iotdb.db.queryengine.execution.driver.DriverContext;
@@ -79,6 +80,7 @@ import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Supplier;
 import java.util.stream.Collectors;
 
+import static com.google.common.base.Preconditions.checkArgument;
 import static 
org.apache.iotdb.db.queryengine.common.DataNodeEndPoints.isSameNode;
 import static 
org.apache.iotdb.db.queryengine.common.FragmentInstanceId.createFullId;
 import static 
org.apache.iotdb.db.queryengine.metric.DataExchangeCostMetricSet.GET_DATA_BLOCK_TASK_SERVER;
@@ -176,6 +178,48 @@ public class MPPDataExchangeManager implements 
IMPPDataExchangeManager {
         }
         // index of the channel must be a SinkChannel
         SinkChannel sinkChannel = (SinkChannel) 
(sinkHandle.getChannel(req.getIndex()));
+        if (req.isSetOffset()) {
+          int remainingPayloadSize =
+              
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangeMaxPayloadSizeInBytes();
+          int offset = req.getOffset();
+          for (int i = req.getStartSequenceId(); i < req.getEndSequenceId(); 
i++) {
+            try {
+              ByteBuffer serializedTsBlock = 
sinkChannel.getSerializedTsBlock(i);
+              int blockOffset = i == req.getStartSequenceId() ? offset : 0;
+              int serializedTsBlockSize = serializedTsBlock.remaining();
+              checkArgument(
+                  blockOffset >= 0 && blockOffset <= serializedTsBlockSize,
+                  DataNodeQueryMessages
+                      
.EXCEPTION_INVALID_SERIALIZED_TSBLOCK_FRAGMENT_OFFSET_ARG_FOR_BLOCK_SIZE_ARG_53BC0284,
+                  blockOffset,
+                  serializedTsBlockSize);
+              int remainingBlockSize = serializedTsBlockSize - blockOffset;
+              ByteBuffer fragment = serializedTsBlock;
+              fragment.position(blockOffset);
+              if (remainingBlockSize <= remainingPayloadSize) {
+                fragment.limit(blockOffset + remainingBlockSize);
+                resp.addToTsBlocks(fragment.slice());
+                remainingPayloadSize -= remainingBlockSize;
+                if (remainingPayloadSize == 0) {
+                  break;
+                }
+              } else {
+                fragment.limit(blockOffset + remainingPayloadSize);
+                resp.addToTsBlocks(fragment.slice());
+                resp.setOffset(blockOffset + remainingPayloadSize);
+                if (blockOffset == 0) {
+                  resp.setTotalLength(serializedTsBlockSize);
+                }
+                break;
+              }
+            } catch (GetTsBlockFromClosedOrAbortedChannelException e) {
+              return new TGetDataBlockResponse(new ArrayList<>());
+            } catch (IllegalArgumentException | IllegalStateException | 
IOException e) {
+              throw new TException(e);
+            }
+          }
+          return resp;
+        }
         for (int i = req.getStartSequenceId(); i < req.getEndSequenceId(); 
i++) {
           try {
             ByteBuffer serializedTsBlock = sinkChannel.getSerializedTsBlock(i);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/sink/SinkChannel.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/sink/SinkChannel.java
index daaff9fdbcd..d24a1552a5b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/sink/SinkChannel.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/sink/SinkChannel.java
@@ -44,7 +44,6 @@ import org.apache.tsfile.common.conf.TSFileDescriptor;
 import org.apache.tsfile.external.commons.lang3.Validate;
 import org.apache.tsfile.read.common.block.TsBlock;
 import org.apache.tsfile.read.common.block.column.TsBlockSerde;
-import org.apache.tsfile.utils.Pair;
 import org.apache.tsfile.utils.RamUsageEstimator;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -67,6 +66,19 @@ import static 
org.apache.iotdb.db.queryengine.metric.DataExchangeCountMetricSet.
 
 public class SinkChannel implements ISinkChannel {
 
+  private static class TsBlockInfo {
+    private TsBlock tsBlock;
+    private ByteBuffer serializedTsBlock;
+    private final long tsBlockSize;
+
+    // private Class, no need to do access control
+    private TsBlockInfo(TsBlock tsBlock, ByteBuffer serializedTsBlock, long 
tsBlockSize) {
+      this.tsBlock = tsBlock;
+      this.serializedTsBlock = serializedTsBlock;
+      this.tsBlockSize = tsBlockSize;
+    }
+  }
+
   private static final Logger LOGGER = 
LoggerFactory.getLogger(SinkChannel.class);
 
   public static final int MAX_ATTEMPT_TIMES = 3;
@@ -93,8 +105,7 @@ public class SinkChannel implements ISinkChannel {
   // Use LinkedHashMap to meet 2 needs,
   //   1. Predictable iteration order so that removing buffered TsBlocks can 
be efficient.
   //   2. Fast lookup.
-  private final LinkedHashMap<Integer, Pair<TsBlock, Long>> 
sequenceIdToTsBlock =
-      new LinkedHashMap<>();
+  private final LinkedHashMap<Integer, TsBlockInfo> sequenceIdToTsBlock = new 
LinkedHashMap<>();
 
   // size for current TsBlock to reserve and free
   private long currentTsBlockSize;
@@ -274,7 +285,7 @@ public class SinkChannel implements ISinkChannel {
       blocked = reserveResult.getFuture();
       bufferRetainedSizeInBytes += reserveResult.getReservedBytes();
 
-      sequenceIdToTsBlock.put(nextSequenceId, new Pair<>(tsBlock, 
currentTsBlockSize));
+      sequenceIdToTsBlock.put(nextSequenceId, new TsBlockInfo(tsBlock, null, 
currentTsBlockSize));
       nextSequenceId += 1;
       currentTsBlockSize = reserveResult.getReservedBytes();
 
@@ -407,8 +418,20 @@ public class SinkChannel implements ISinkChannel {
       throw new GetTsBlockFromClosedOrAbortedChannelException(
           DataNodeQueryMessages.SINKCHANNEL_IS_ABORTED_OR_CLOSED);
     }
-    Pair<TsBlock, Long> pair = sequenceIdToTsBlock.get(sequenceId);
-    if (pair == null || pair.left == null) {
+    TsBlockInfo tsBlockInfo = sequenceIdToTsBlock.get(sequenceId);
+    if (tsBlockInfo == null) {
+      LOGGER.warn(
+          DataNodeQueryMessages.THE_TSBLOCK_DOESNT_EXIST_SEQUENCE_ID_REMAINING,
+          sequenceId,
+          sequenceIdToTsBlock.entrySet());
+      throw new IllegalStateException(
+          DataNodeQueryMessages.THE_DATA_BLOCK_DOESN_T_EXIST_SEQUENCE_ID + 
sequenceId);
+    }
+    ByteBuffer serializedTsBlock = tsBlockInfo.serializedTsBlock;
+    if (serializedTsBlock != null) {
+      return serializedTsBlock.duplicate();
+    }
+    if (tsBlockInfo.tsBlock == null) {
       LOGGER.warn(
           DataNodeQueryMessages.THE_TSBLOCK_DOESNT_EXIST_SEQUENCE_ID_REMAINING,
           sequenceId,
@@ -416,7 +439,10 @@ public class SinkChannel implements ISinkChannel {
       throw new IllegalStateException(
           DataNodeQueryMessages.THE_DATA_BLOCK_DOESN_T_EXIST_SEQUENCE_ID + 
sequenceId);
     }
-    return serde.serialize(pair.left);
+    serializedTsBlock = serde.serialize(tsBlockInfo.tsBlock);
+    tsBlockInfo.serializedTsBlock = serializedTsBlock;
+    tsBlockInfo.tsBlock = null;
+    return serializedTsBlock.duplicate();
   }
 
   public void acknowledgeTsBlock(int startSequenceId, int endSequenceId) {
@@ -425,10 +451,9 @@ public class SinkChannel implements ISinkChannel {
       if (aborted || closed) {
         return;
       }
-      Iterator<Entry<Integer, Pair<TsBlock, Long>>> iterator =
-          sequenceIdToTsBlock.entrySet().iterator();
+      Iterator<Entry<Integer, TsBlockInfo>> iterator = 
sequenceIdToTsBlock.entrySet().iterator();
       while (iterator.hasNext()) {
-        Entry<Integer, Pair<TsBlock, Long>> entry = iterator.next();
+        Entry<Integer, TsBlockInfo> entry = iterator.next();
         if (entry.getKey() < startSequenceId) {
           continue;
         }
@@ -436,8 +461,8 @@ public class SinkChannel implements ISinkChannel {
           break;
         }
 
-        freedBytes += entry.getValue().right;
-        bufferRetainedSizeInBytes -= entry.getValue().right;
+        freedBytes += entry.getValue().tsBlockSize;
+        bufferRetainedSizeInBytes -= entry.getValue().tsBlockSize;
         iterator.remove();
         if (LOGGER.isDebugEnabled()) {
           LOGGER.debug(DataNodeQueryMessages.ACK_TSBLOCK, entry.getKey());
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
index 67a5defb09a..aed3f9533fb 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/execution/exchange/source/SourceHandle.java
@@ -40,6 +40,7 @@ import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstanceId;
 import org.apache.iotdb.mpp.rpc.thrift.TGetDataBlockRequest;
 import org.apache.iotdb.mpp.rpc.thrift.TGetDataBlockResponse;
 
+import com.google.common.base.Preconditions;
 import com.google.common.util.concurrent.ListenableFuture;
 import com.google.common.util.concurrent.SettableFuture;
 import org.apache.thrift.TException;
@@ -47,6 +48,7 @@ import org.apache.tsfile.external.commons.lang3.Validate;
 import org.apache.tsfile.read.common.block.TsBlock;
 import org.apache.tsfile.read.common.block.column.TsBlockSerde;
 import org.apache.tsfile.utils.Pair;
+import org.apache.tsfile.utils.PublicBAOS;
 import org.apache.tsfile.utils.RamUsageEstimator;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -59,6 +61,7 @@ import java.util.Map;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 
+import static com.google.common.base.Preconditions.checkArgument;
 import static 
com.google.common.util.concurrent.Futures.nonCancellationPropagating;
 import static 
org.apache.iotdb.db.queryengine.execution.exchange.MPPDataExchangeManager.createFullIdFrom;
 import static 
org.apache.iotdb.db.queryengine.metric.DataExchangeCostMetricSet.GET_DATA_BLOCK_TASK_CALLER;
@@ -640,6 +643,8 @@ public class SourceHandle implements ISourceHandle {
                 startSequenceId,
                 endSequenceId,
                 indexOfUpstreamSinkHandle);
+        DataBlockFetchProgress fetchProgress =
+            new DataBlockFetchProgress(startSequenceId, endSequenceId);
         int attempt = 0;
         while (attempt < MAX_ATTEMPT_TIMES) {
           attempt += 1;
@@ -648,7 +653,10 @@ public class SourceHandle implements ISourceHandle {
           boolean transferAttemptRecorded = false;
           try (SyncDataNodeMPPDataExchangeServiceClient client =
               
mppDataExchangeServiceClientManager.borrowClient(remoteEndpoint)) {
-            TGetDataBlockResponse resp = client.getDataBlock(req);
+            TGetDataBlockResponse resp = getDataBlockWithFragments(client, 
req, fetchProgress);
+            if (resp == null) {
+              return;
+            }
             int tsBlockNum = resp.getTsBlocks().size();
             if (tsBlockNum != endSequenceId - startSequenceId) {
               recordTransferAttempt(
@@ -709,7 +717,25 @@ public class SourceHandle implements ISourceHandle {
               return;
             }
             break;
-          } catch (Throwable e) {
+          } catch (IllegalArgumentException | IllegalStateException e) {
+            if (!transferAttemptRecorded) {
+              recordTransferAttempt(
+                  false, 
UserDataTransferErrorCode.UNEXPECTED_RESPONSE_SIZE.name(), e);
+            }
+            fail(e);
+            return;
+          } catch (Error e) {
+            if (!transferAttemptRecorded) {
+              recordTransferAttempt(
+                  false,
+                  e instanceof OutOfMemoryError
+                      ? UserDataTransferErrorCode.OUT_OF_MEMORY.name()
+                      : null,
+                  e);
+            }
+            fail(e);
+            throw e;
+          } catch (Exception e) {
 
             if (!transferAttemptRecorded) {
               recordTransferAttempt(false, null, e);
@@ -767,6 +793,134 @@ public class SourceHandle implements ISourceHandle {
         sourceHandleListener.onFailure(SourceHandle.this, t);
       }
     }
+
+    private TGetDataBlockResponse getDataBlockWithFragments(
+        SyncDataNodeMPPDataExchangeServiceClient client,
+        TGetDataBlockRequest request,
+        DataBlockFetchProgress fetchProgress)
+        throws TException {
+      while (!fetchProgress.isFinished()) {
+        synchronized (SourceHandle.this) {
+          if (aborted || closed) {
+            fetchProgress.discard();
+            return null;
+          }
+        }
+        TGetDataBlockRequest fragmentRequest = request.deepCopy();
+        fragmentRequest.setStartSequenceId(fetchProgress.nextSequenceId);
+        fragmentRequest.setOffset(fetchProgress.offset);
+        TGetDataBlockResponse response = client.getDataBlock(fragmentRequest);
+        synchronized (SourceHandle.this) {
+          if (aborted || closed) {
+            fetchProgress.discard();
+            return null;
+          }
+          if (response.getTsBlocks().isEmpty()) {
+            return response;
+          }
+          fetchProgress.addResponse(response);
+        }
+      }
+      return new TGetDataBlockResponse(fetchProgress.tsBlocks);
+    }
+
+    private class DataBlockFetchProgress {
+      private final int endSequenceId;
+      private final List<ByteBuffer> tsBlocks;
+      private int nextSequenceId;
+      private int offset;
+      private PublicBAOS partialTsBlock;
+      private int partialTsBlockTotalLength;
+
+      private DataBlockFetchProgress(int startSequenceId, int endSequenceId) {
+        this.nextSequenceId = startSequenceId;
+        this.endSequenceId = endSequenceId;
+        this.tsBlocks = new ArrayList<>(endSequenceId - startSequenceId);
+      }
+
+      private boolean isFinished() {
+        return nextSequenceId == endSequenceId && partialTsBlock == null;
+      }
+
+      private void discard() {
+        tsBlocks.clear();
+        partialTsBlock = null;
+        partialTsBlockTotalLength = 0;
+        offset = 0;
+      }
+
+      private void addResponse(TGetDataBlockResponse response) {
+        List<ByteBuffer> responseBlocks = response.getTsBlocks();
+        boolean lastBlockIsFragment = response.isSetOffset();
+        int blockIndex = 0;
+
+        if (partialTsBlock != null) {
+          appendFragment(responseBlocks.get(blockIndex++));
+          if (lastBlockIsFragment && blockIndex == responseBlocks.size()) {
+            updateOffset(response.getOffset());
+            return;
+          }
+          Preconditions.checkState(
+              partialTsBlock.size() == partialTsBlockTotalLength,
+              DataNodeQueryMessages
+                  
.EXCEPTION_ACCUMULATED_TSBLOCK_FRAGMENT_LENGTH_ARG_DOES_NOT_MATCH_TOTALLENGTH_ARG_1B784303,
+              partialTsBlock.size(),
+              partialTsBlockTotalLength);
+          tsBlocks.add(ByteBuffer.wrap(partialTsBlock.getBuf(), 0, 
partialTsBlock.size()));
+          partialTsBlock = null;
+          partialTsBlockTotalLength = 0;
+          offset = 0;
+          nextSequenceId++;
+        }
+
+        int lastCompleteBlockIndex =
+            lastBlockIsFragment ? responseBlocks.size() - 1 : 
responseBlocks.size();
+        while (blockIndex < lastCompleteBlockIndex) {
+          tsBlocks.add(responseBlocks.get(blockIndex++));
+          nextSequenceId++;
+        }
+
+        if (lastBlockIsFragment) {
+          checkArgument(
+              response.isSetTotalLength(),
+              DataNodeQueryMessages
+                  
.EXCEPTION_THE_FIRST_FRAGMENTED_DATA_BLOCK_RESPONSE_MUST_INCLUDE_TOTALLENGTH_C5C79BC2);
+          partialTsBlockTotalLength = response.getTotalLength();
+          partialTsBlock = new PublicBAOS(partialTsBlockTotalLength);
+          appendFragment(responseBlocks.get(blockIndex));
+          updateOffset(response.getOffset());
+        }
+
+        Preconditions.checkState(
+            nextSequenceId <= endSequenceId,
+            DataNodeQueryMessages
+                
.EXCEPTION_NEXT_SEQUENCE_ID_ARG_EXCEEDS_REQUESTED_END_SEQUENCE_ID_ARG_30B1726E,
+            nextSequenceId,
+            endSequenceId);
+        Preconditions.checkState(
+            lastBlockIsFragment || nextSequenceId != endSequenceId || 
partialTsBlock == null,
+            DataNodeQueryMessages
+                
.EXCEPTION_A_COMPLETED_DATA_BLOCK_RESPONSE_RANGE_MUST_NOT_RETAIN_A_PARTIAL_TSBLOCK_85E5C287);
+      }
+
+      private void appendFragment(ByteBuffer fragment) {
+        checkArgument(
+            fragment.hasRemaining(),
+            
DataNodeQueryMessages.EXCEPTION_TSBLOCK_FRAGMENT_MUST_NOT_BE_EMPTY_C7D19863);
+        partialTsBlock.writeBytes(fragment.array());
+      }
+
+      private void updateOffset(int nextOffset) {
+        checkArgument(
+            nextOffset > offset && nextOffset == partialTsBlock.size(),
+            DataNodeQueryMessages
+                
.EXCEPTION_NEXT_FRAGMENT_OFFSET_ARG_MUST_BE_GREATER_THAN_CURRENT_OFFSET_ARG_AND_MATCH_ACCUMULATED_FRAGMENT_LENGTH_ARG_ECC31047,
+            nextOffset,
+            offset,
+            partialTsBlock.size());
+        offset = nextOffset;
+      }
+    }
   }
 
   class SendAcknowledgeDataBlockEventTask implements Runnable {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
index be8b29a064b..b0d7aef7394 100755
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/conf/PropertiesTest.java
@@ -295,6 +295,34 @@ public class PropertiesTest {
     }
   }
 
+  @Test
+  public void testMppDataExchangeMaxPayloadSizeHotReload() throws Exception {
+    IoTDBDescriptor descriptor = IoTDBDescriptor.getInstance();
+    int originalPayloadSize = 
descriptor.getConfig().getMppDataExchangeMaxPayloadSizeInBytes();
+    try {
+      TrimProperties properties = new TrimProperties();
+
+      properties.setProperty("mpp_data_exchange_max_payload_size_in_bytes", 
"0");
+      descriptor.loadHotModifiedProps(properties);
+      Assert.assertEquals(
+          4 * 1024 * 1024, 
descriptor.getConfig().getMppDataExchangeMaxPayloadSizeInBytes());
+
+      properties.setProperty("mpp_data_exchange_max_payload_size_in_bytes", 
"1");
+      descriptor.loadHotModifiedProps(properties);
+      Assert.assertEquals(
+          128 * 1024, 
descriptor.getConfig().getMppDataExchangeMaxPayloadSizeInBytes());
+
+      int maximumPayloadSize = descriptor.getConfig().getThriftMaxFrameSize() 
- 1024;
+      properties.setProperty(
+          "mpp_data_exchange_max_payload_size_in_bytes", 
Integer.toString(Integer.MAX_VALUE));
+      descriptor.loadHotModifiedProps(properties);
+      Assert.assertEquals(
+          maximumPayloadSize, 
descriptor.getConfig().getMppDataExchangeMaxPayloadSizeInBytes());
+    } finally {
+      
descriptor.getConfig().setMppDataExchangeMaxPayloadSizeInBytes(originalPayloadSize);
+    }
+  }
+
   @Test
   public void PropertiesWithSpace() {
     IoTDBDescriptor descriptor = IoTDBDescriptor.getInstance();
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SinkChannelTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SinkChannelTest.java
index f59dd9925aa..f22725c0fea 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SinkChannelTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SinkChannelTest.java
@@ -37,6 +37,7 @@ import org.apache.iotdb.mpp.rpc.thrift.TNewDataBlockEvent;
 import org.apache.thrift.TException;
 import org.apache.tsfile.common.conf.TSFileDescriptor;
 import org.apache.tsfile.read.common.block.TsBlock;
+import org.apache.tsfile.read.common.block.column.TsBlockSerde;
 import org.junit.Assert;
 import org.junit.Test;
 import org.mockito.Mockito;
@@ -87,6 +88,7 @@ public class SinkChannelTest {
     SinkListener mockSinkListener = Mockito.mock(SinkListener.class);
     // Construct several mock TsBlock(s).
     List<TsBlock> mockTsBlocks = Utils.createMockTsBlocks(numOfMockTsBlock, 
mockTsBlockSize);
+    TsBlockSerde mockTsBlockSerde = 
Utils.createMockTsBlockSerde(mockTsBlockSize);
 
     // Construct SinkChannel.
     SinkChannel sinkChannel =
@@ -98,7 +100,7 @@ public class SinkChannelTest {
             localFragmentInstanceId,
             mockLocalMemoryManager,
             Executors.newSingleThreadExecutor(),
-            Utils.createMockTsBlockSerde(mockTsBlockSize),
+            mockTsBlockSerde,
             mockSinkListener,
             mockClientManager);
     sinkChannel.open();
@@ -145,6 +147,8 @@ public class SinkChannelTest {
     for (int i = 0; i < numOfMockTsBlock; i++) {
       try {
         sinkChannel.getSerializedTsBlock(i);
+        sinkChannel.getSerializedTsBlock(i);
+        Mockito.verify(mockTsBlockSerde, 
Mockito.times(1)).serialize(mockTsBlocks.get(i));
       } catch (IOException e) {
         e.printStackTrace();
         Assert.fail();
@@ -166,6 +170,16 @@ public class SinkChannelTest {
     Assert.assertTrue(sinkChannel.isFinished());
     Assert.assertFalse(sinkChannel.isAborted());
     Assert.assertEquals(mockTsBlockSize, 
sinkChannel.getBufferRetainedSizeInBytes());
+    for (int i = 0; i < numOfMockTsBlock; i++) {
+      try {
+        sinkChannel.getSerializedTsBlock(i);
+        Assert.fail("The acknowledged serialized TsBlock should have been 
released");
+      } catch (IllegalStateException expected) {
+        // Both the original entry and serialized cache must be removed by 
acknowledgement.
+      } catch (IOException e) {
+        Assert.fail(e.getMessage());
+      }
+    }
     Mockito.verify(mockMemoryPool, Mockito.timeout(10_0000).times(1))
         .free(
             queryId,
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
index 03220a3d10b..af96360b12e 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/SourceHandleTest.java
@@ -46,8 +46,11 @@ import org.mockito.Mockito;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.List;
+import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
@@ -70,6 +73,181 @@ public class SourceHandleTest {
     
IoTDBDescriptor.getInstance().getMemoryConfig().setMaxBytesPerFragmentInstance(maxBytesPerFI);
   }
 
+  @Test
+  public void testCloseStopsFragmentFetchAfterInFlightRpc() throws Exception {
+    testCancellationStopsFragmentFetchAfterInFlightRpc(false);
+  }
+
+  @Test
+  public void testAbortStopsFragmentFetchAfterInFlightRpc() throws Exception {
+    testCancellationStopsFragmentFetchAfterInFlightRpc(true);
+  }
+
+  private void testCancellationStopsFragmentFetchAfterInFlightRpc(boolean 
abort) throws Exception {
+    final String queryId = "q0";
+    final TEndPoint remoteEndpoint =
+        new TEndPoint("remote", 
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangePort());
+    final TFragmentInstanceId remoteFragmentInstanceId = new 
TFragmentInstanceId(queryId, 1, "0");
+    final TFragmentInstanceId localFragmentInstanceId = new 
TFragmentInstanceId(queryId, 0, "0");
+    final LocalMemoryManager localMemoryManager = 
Mockito.mock(LocalMemoryManager.class);
+    final MemoryPool memoryPool = Utils.createMockNonBlockedMemoryPool();
+    Mockito.when(localMemoryManager.getQueryPool()).thenReturn(memoryPool);
+    final SourceHandleListener sourceHandleListener = 
Mockito.mock(SourceHandleListener.class);
+    final IClientManager<TEndPoint, SyncDataNodeMPPDataExchangeServiceClient> 
clientManager =
+        Mockito.mock(IClientManager.class);
+    final SyncDataNodeMPPDataExchangeServiceClient client =
+        Mockito.mock(SyncDataNodeMPPDataExchangeServiceClient.class);
+    final CountDownLatch rpcStarted = new CountDownLatch(1);
+    final CountDownLatch returnResponse = new CountDownLatch(1);
+    
Mockito.when(clientManager.borrowClient(remoteEndpoint)).thenReturn(client);
+    Mockito.doAnswer(
+            invocation -> {
+              rpcStarted.countDown();
+              Assert.assertTrue(returnResponse.await(10, TimeUnit.SECONDS));
+              return new TGetDataBlockResponse(List.of(ByteBuffer.wrap(new 
byte[] {1, 2})))
+                  .setOffset(2)
+                  .setTotalLength(6);
+            })
+        .when(client)
+        .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+
+    final SourceHandle sourceHandle =
+        new SourceHandle(
+            remoteEndpoint,
+            remoteFragmentInstanceId,
+            localFragmentInstanceId,
+            "exchange_0",
+            0,
+            localMemoryManager,
+            Executors.newSingleThreadExecutor(),
+            Utils.createMockTsBlockSerde(MOCK_TSBLOCK_SIZE),
+            sourceHandleListener,
+            clientManager);
+    sourceHandle.isBlocked();
+    sourceHandle.updatePendingDataBlockInfo(0, List.of(MOCK_TSBLOCK_SIZE));
+    Assert.assertTrue(rpcStarted.await(10, TimeUnit.SECONDS));
+
+    if (abort) {
+      sourceHandle.abort();
+    } else {
+      sourceHandle.close();
+    }
+    returnResponse.countDown();
+
+    Mockito.verify(client, Mockito.timeout(10_000).times(1))
+        .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    Mockito.verify(client, Mockito.after(500).times(1))
+        .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    Mockito.verify(sourceHandleListener, Mockito.never())
+        .onFailure(Mockito.eq(sourceHandle), Mockito.any());
+    Assert.assertEquals(0L, sourceHandle.getBufferRetainedSizeInBytes());
+  }
+
+  @Test
+  public void testFragmentAllocationErrorFailsWithoutRetry() throws Exception {
+    final String queryId = "q0";
+    final TEndPoint remoteEndpoint =
+        new TEndPoint("remote", 
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangePort());
+    final TFragmentInstanceId remoteFragmentInstanceId = new 
TFragmentInstanceId(queryId, 1, "0");
+    final TFragmentInstanceId localFragmentInstanceId = new 
TFragmentInstanceId(queryId, 0, "0");
+    final LocalMemoryManager localMemoryManager = 
Mockito.mock(LocalMemoryManager.class);
+    final MemoryPool memoryPool = Utils.createMockNonBlockedMemoryPool();
+    Mockito.when(localMemoryManager.getQueryPool()).thenReturn(memoryPool);
+    final SourceHandleListener sourceHandleListener = 
Mockito.mock(SourceHandleListener.class);
+    final IClientManager<TEndPoint, SyncDataNodeMPPDataExchangeServiceClient> 
clientManager =
+        Mockito.mock(IClientManager.class);
+    final SyncDataNodeMPPDataExchangeServiceClient client =
+        Mockito.mock(SyncDataNodeMPPDataExchangeServiceClient.class);
+    
Mockito.when(clientManager.borrowClient(remoteEndpoint)).thenReturn(client);
+    Mockito.when(client.getDataBlock(Mockito.any(TGetDataBlockRequest.class)))
+        .thenReturn(
+            new TGetDataBlockResponse(List.of(ByteBuffer.wrap(new byte[] {1})))
+                .setOffset(1)
+                .setTotalLength(Integer.MAX_VALUE));
+
+    final SourceHandle sourceHandle =
+        new SourceHandle(
+            remoteEndpoint,
+            remoteFragmentInstanceId,
+            localFragmentInstanceId,
+            "exchange_0",
+            0,
+            localMemoryManager,
+            Executors.newSingleThreadExecutor(),
+            Utils.createMockTsBlockSerde(MOCK_TSBLOCK_SIZE),
+            sourceHandleListener,
+            clientManager);
+    sourceHandle.isBlocked();
+    sourceHandle.updatePendingDataBlockInfo(0, List.of(MOCK_TSBLOCK_SIZE));
+
+    Mockito.verify(sourceHandleListener, Mockito.timeout(10_000))
+        .onFailure(Mockito.eq(sourceHandle), 
Mockito.any(OutOfMemoryError.class));
+    Mockito.verify(client, 
Mockito.times(1)).getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    Assert.assertEquals(0L, sourceHandle.getBufferRetainedSizeInBytes());
+    sourceHandle.abort();
+  }
+
+  @Test
+  public void testFragmentAssemblyFailsWhenTotalLengthDoesNotMatch() {
+    final String queryId = "q0";
+    final TEndPoint remoteEndpoint =
+        new TEndPoint("remote", 
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangePort());
+    final TFragmentInstanceId remoteFragmentInstanceId = new 
TFragmentInstanceId(queryId, 1, "0");
+    final TFragmentInstanceId localFragmentInstanceId = new 
TFragmentInstanceId(queryId, 0, "0");
+
+    LocalMemoryManager localMemoryManager = 
Mockito.mock(LocalMemoryManager.class);
+    MemoryPool memoryPool = Utils.createMockNonBlockedMemoryPool();
+    Mockito.when(localMemoryManager.getQueryPool()).thenReturn(memoryPool);
+    SourceHandleListener sourceHandleListener = 
Mockito.mock(SourceHandleListener.class);
+    TsBlockSerde serde = Utils.createMockTsBlockSerde(MOCK_TSBLOCK_SIZE);
+    IClientManager<TEndPoint, SyncDataNodeMPPDataExchangeServiceClient> 
clientManager =
+        Mockito.mock(IClientManager.class);
+    SyncDataNodeMPPDataExchangeServiceClient client =
+        Mockito.mock(SyncDataNodeMPPDataExchangeServiceClient.class);
+    AtomicInteger rpcCount = new AtomicInteger();
+    try {
+      
Mockito.when(clientManager.borrowClient(remoteEndpoint)).thenReturn(client);
+      Mockito.doAnswer(
+              invocation -> {
+                if (rpcCount.getAndIncrement() == 0) {
+                  return new TGetDataBlockResponse(List.of(ByteBuffer.wrap(new 
byte[] {1, 2})))
+                      .setOffset(2)
+                      .setTotalLength(5);
+                }
+                return new TGetDataBlockResponse(List.of(ByteBuffer.wrap(new 
byte[] {3, 4})));
+              })
+          .when(client)
+          .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    } catch (ClientManagerException | TException e) {
+      Assert.fail(e.getMessage());
+    }
+
+    SourceHandle sourceHandle =
+        new SourceHandle(
+            remoteEndpoint,
+            remoteFragmentInstanceId,
+            localFragmentInstanceId,
+            "exchange_0",
+            0,
+            localMemoryManager,
+            Executors.newSingleThreadExecutor(),
+            serde,
+            sourceHandleListener,
+            clientManager);
+    Assert.assertFalse(sourceHandle.isBlocked().isDone());
+    sourceHandle.updatePendingDataBlockInfo(0, List.of(MOCK_TSBLOCK_SIZE));
+    try {
+      Mockito.verify(client, Mockito.timeout(10_000).times(2))
+          .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
+    } catch (TException e) {
+      Assert.fail(e.getMessage());
+    }
+    Mockito.verify(sourceHandleListener, Mockito.timeout(10_000))
+        .onFailure(Mockito.eq(sourceHandle), 
Mockito.any(IllegalStateException.class));
+    Mockito.verify(serde, 
Mockito.never()).deserialize(Mockito.any(ByteBuffer.class));
+    sourceHandle.abort();
+  }
+
   @Test
   public void testNonBlockedOneTimeReceive() {
     final String queryId = "q0";
@@ -634,11 +812,14 @@ public class SourceHandleTest {
       Mockito.doAnswer(
               invocation -> {
                 final TGetDataBlockRequest request = invocation.getArgument(0);
+                if (request.getStartSequenceId() > 0) {
+                  throw new TException("mock RPC failure");
+                }
                 final List<ByteBuffer> shortResponse = new ArrayList<>();
                 for (int i = 0;
                     i < request.getEndSequenceId() - 
request.getStartSequenceId() - 1;
                     i++) {
-                  shortResponse.add(ByteBuffer.allocate(0));
+                  shortResponse.add(ByteBuffer.allocate(1));
                 }
                 return new TGetDataBlockResponse(shortResponse);
               })
@@ -670,7 +851,7 @@ public class SourceHandleTest {
             .collect(Collectors.toList()));
 
     try {
-      Mockito.verify(mockClient, 
Mockito.timeout(10_000).times(SourceHandle.MAX_ATTEMPT_TIMES))
+      Mockito.verify(mockClient, 
Mockito.timeout(10_000).times(SourceHandle.MAX_ATTEMPT_TIMES + 1))
           .getDataBlock(Mockito.any(TGetDataBlockRequest.class));
     } catch (TException e) {
       Assert.fail(e.getMessage());
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/Utils.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/Utils.java
index b09498ad949..1a701b9eded 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/Utils.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/exchange/Utils.java
@@ -28,6 +28,7 @@ import 
org.apache.tsfile.read.common.block.column.TsBlockSerde;
 import org.mockito.Mockito;
 import org.mockito.stubbing.Answer;
 
+import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.List;
@@ -155,6 +156,12 @@ public class Utils {
     TsBlock mockTsBlock = Mockito.mock(TsBlock.class);
     
Mockito.when(mockTsBlock.getRetainedSizeInBytes()).thenReturn(mockTsBlockSize);
     Mockito.when(mockTsBlock.getSizeInBytes()).thenReturn(mockTsBlockSize);
+    try {
+      Mockito.when(mockTsBlockSerde.serialize(Mockito.any(TsBlock.class)))
+          .thenReturn(ByteBuffer.allocate(Math.toIntExact(mockTsBlockSize)));
+    } catch (IOException e) {
+      throw new AssertionError(e);
+    }
     Mockito.when(mockTsBlockSerde.deserialize(Mockito.any(ByteBuffer.class)))
         .thenReturn(mockTsBlock);
     return mockTsBlockSerde;
diff --git 
a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
 
b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
index 5079abe5100..2c6186f6b4d 100644
--- 
a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
+++ 
b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template
@@ -1194,6 +1194,13 @@ mpp_data_exchange_max_pool_size=10
 # Datatype: int
 mpp_data_exchange_keep_alive_time_in_ms=1000
 
+# The maximum payload size of one MPP data exchange RPC response.
+# <=0 use default value
+# The effective value is between 131072(128KB) and dn_thrift_max_frame_size 
minus 1024 bytes reserved for the RPC response envelope.
+# effectiveMode: hot_reload
+# Datatype: int, Unit: byte
+mpp_data_exchange_max_payload_size_in_bytes=4194304
+
 # The max execution time of a DriverTask
 # effectiveMode: restart
 # Datatype: int, Unit: ms
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
index 2bf187999ed..18e9b5f9337 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/audit/UserDataTransferErrorCode.java
@@ -23,5 +23,6 @@ public enum UserDataTransferErrorCode {
   EMPTY_RESPONSE,
   UNEXPECTED_RESPONSE_SIZE,
   RECEIVER_CLOSED,
-  REMOTE_REJECTED
+  REMOTE_REJECTED,
+  OUT_OF_MEMORY
 }
diff --git a/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift 
b/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift
index c20c1099888..73bd95942f1 100644
--- a/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift
+++ b/iotdb-protocol/thrift-datanode/src/main/thrift/datanode.thrift
@@ -90,10 +90,16 @@ struct TGetDataBlockRequest {
   3: required i32 endSequenceId
   // Index of upstream SinkChannel
   4: required i32 index
+  // Optional byte range for fetching one serialized TsBlock in fragments.
+  5: optional i32 offset
 }
 
 struct TGetDataBlockResponse {
   1: required list<binary> tsBlocks
+  // The start offset of the next fragment. It is set only when the last 
element in tsBlocks is a fragment.
+  2: optional i32 offset
+  // Total serialized length of the TsBlock when the response starts its first 
fragment.
+  3: optional i32 totalLength
 }
 
 struct TAcknowledgeDataBlockEvent {

Reply via email to