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

Wei-hao-Li pushed a commit to branch mppEx
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 62fd96b325217c61b08dd59a04132823e5268a0c
Author: Weihao Li <[email protected]>
AuthorDate: Mon Sep 14 16:45:25 2026 +0800

    optimize some
    
    Signed-off-by: Weihao Li <[email protected]>
---
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |  3 +--
 .../execution/exchange/MPPDataExchangeManager.java | 30 ++++++++++++++--------
 .../execution/exchange/source/SourceHandle.java    |  7 ++---
 3 files changed, 24 insertions(+), 16 deletions(-)

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 7d7f3652d38..ff141a8ed85 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
@@ -3427,8 +3427,7 @@ public class IoTDBConfig {
     return mppDataExchangeMaxPayloadSizeInBytes;
   }
 
-  public void setMppDataExchangeMaxPayloadSizeInBytes(
-      int mppDataExchangeMaxPayloadSizeInBytes) {
+  public void setMppDataExchangeMaxPayloadSizeInBytes(int 
mppDataExchangeMaxPayloadSizeInBytes) {
     this.mppDataExchangeMaxPayloadSizeInBytes =
         Math.max(1, Math.min(mppDataExchangeMaxPayloadSizeInBytes, 
thriftMaxFrameSize - 1024));
   }
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 ea3dce68281..e60a7b3d89f 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
@@ -179,27 +179,35 @@ public class MPPDataExchangeManager implements 
IMPPDataExchangeManager {
         SinkChannel sinkChannel = (SinkChannel) 
(sinkHandle.getChannel(req.getIndex()));
         if (req.isSetOffset()) {
           int remainingPayloadSize =
-              IoTDBDescriptor.getInstance()
-                  .getConfig()
-                  .getMppDataExchangeMaxPayloadSizeInBytes();
+              
IoTDBDescriptor.getInstance().getConfig().getMppDataExchangeMaxPayloadSizeInBytes();
           int offset = req.getOffset();
           for (int i = req.getStartSequenceId(); i < req.getEndSequenceId(); 
i++) {
             try {
-              ByteBuffer serializedTsBlock = 
sinkChannel.getSerializedTsBlock(i);
+              ByteBuffer serializedTsBlock = 
sinkChannel.getSerializedTsBlock(i).asReadOnlyBuffer();
               int blockOffset = i == req.getStartSequenceId() ? offset : 0;
-              int remainingBlockSize = serializedTsBlock.remaining() - 
blockOffset;
+              int serializedTsBlockSize = serializedTsBlock.remaining();
+              if (blockOffset < 0 || blockOffset > serializedTsBlockSize) {
+                throw new IllegalArgumentException(
+                    String.format(
+                        
DataNodeQueryMessages.EXCEPTION_INVALID_ARG_ARG_2946DBE5,
+                        "serialized TsBlock",
+                        "fragment range"));
+              }
+              int remainingBlockSize = serializedTsBlockSize - blockOffset;
               if (remainingBlockSize <= remainingPayloadSize) {
-                resp.addToTsBlocks(
-                    sinkChannel.getSerializedTsBlockFragment(
-                        i, blockOffset, remainingBlockSize));
+                ByteBuffer fragment = serializedTsBlock.duplicate();
+                fragment.position(blockOffset);
+                fragment.limit(blockOffset + remainingBlockSize);
+                resp.addToTsBlocks(fragment.slice());
                 remainingPayloadSize -= remainingBlockSize;
                 if (remainingPayloadSize == 0) {
                   break;
                 }
               } else {
-                resp.addToTsBlocks(
-                    sinkChannel.getSerializedTsBlockFragment(
-                        i, blockOffset, remainingPayloadSize));
+                ByteBuffer fragment = serializedTsBlock.duplicate();
+                fragment.position(blockOffset);
+                fragment.limit(blockOffset + remainingPayloadSize);
+                resp.addToTsBlocks(fragment.slice());
                 resp.setOffset(blockOffset + remainingPayloadSize);
                 break;
               }
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 db4f1dfd408..3afc800c2af 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
@@ -651,8 +651,7 @@ public class SourceHandle implements ISourceHandle {
           boolean transferAttemptRecorded = false;
           try (SyncDataNodeMPPDataExchangeServiceClient client =
               
mppDataExchangeServiceClientManager.borrowClient(remoteEndpoint)) {
-            TGetDataBlockResponse resp =
-                getDataBlockWithFragments(client, req, fetchProgress);
+            TGetDataBlockResponse resp = getDataBlockWithFragments(client, 
req, fetchProgress);
             int tsBlockNum = resp.getTsBlocks().size();
             if (tsBlockNum != endSequenceId - startSequenceId) {
               recordTransferAttempt(
@@ -838,7 +837,9 @@ public class SourceHandle implements ISourceHandle {
         }
 
         if (nextSequenceId > endSequenceId
-            || (!lastBlockIsFragment && nextSequenceId == endSequenceId && 
partialTsBlock != null)) {
+            || (!lastBlockIsFragment
+                && nextSequenceId == endSequenceId
+                && partialTsBlock != null)) {
           throw new TException(
               
DataNodeQueryMessages.EXCEPTION_UNEXPECTED_DATA_BLOCK_RESPONSE_SIZE_A7DD7E33);
         }

Reply via email to