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); }
