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

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

commit c632b152467e5157c00e5c7d2d540130bbc28ae4
Author: JackieTien97 <[email protected]>
AuthorDate: Wed Jun 8 11:13:16 2022 +0800

    Support LocalSinHandle
---
 .../org/apache/iotdb/db/mpp/plan/Coordinator.java  |  2 +-
 .../db/mpp/plan/planner/LocalExecutionPlanner.java | 48 ++++++++++++++++------
 2 files changed, 36 insertions(+), 14 deletions(-)

diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java
index cf19379297..3b91b52453 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/Coordinator.java
@@ -61,7 +61,7 @@ public class Coordinator {
   private static final String COORDINATOR_SCHEDULED_EXECUTOR_NAME = 
"MPPCoordinatorScheduled";
   private static final int COORDINATOR_SCHEDULED_EXECUTOR_SIZE = 1;
 
-  private static final TEndPoint LOCAL_HOST_DATA_BLOCK_ENDPOINT =
+  public static final TEndPoint LOCAL_HOST_DATA_BLOCK_ENDPOINT =
       new TEndPoint(
           IoTDBDescriptor.getInstance().getConfig().getInternalIp(),
           IoTDBDescriptor.getInstance().getConfig().getDataBlockManagerPort());
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
index 70108b5932..5e0799b22c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LocalExecutionPlanner.java
@@ -18,6 +18,7 @@
  */
 package org.apache.iotdb.db.mpp.plan.planner;
 
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.db.engine.storagegroup.DataRegion;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
@@ -178,6 +179,7 @@ import java.util.stream.Collectors;
 import static com.google.common.base.Preconditions.checkArgument;
 import static java.util.Objects.requireNonNull;
 import static 
org.apache.iotdb.db.mpp.execution.operator.LastQueryUtil.satisfyFilter;
+import static 
org.apache.iotdb.db.mpp.plan.Coordinator.LOCAL_HOST_DATA_BLOCK_ENDPOINT;
 
 /**
  * Used to plan a fragment instance. Currently, we simply change it from 
PlanNode to executable
@@ -1051,17 +1053,24 @@ public class LocalExecutionPlanner {
           context.instanceContext.addOperatorContext(
               context.getNextOperatorId(),
               node.getPlanNodeId(),
-              SeriesScanOperator.class.getSimpleName());
+              ExchangeOperator.class.getSimpleName());
       FragmentInstanceId localInstanceId = context.instanceContext.getId();
       FragmentInstanceId remoteInstanceId = node.getUpstreamInstanceId();
 
+      TEndPoint upstreamEndPoint = node.getUpstreamEndpoint();
       ISourceHandle sourceHandle =
-          DATA_BLOCK_MANAGER.createSourceHandle(
-              localInstanceId.toThrift(),
-              node.getPlanNodeId().getId(),
-              node.getUpstreamEndpoint(),
-              remoteInstanceId.toThrift(),
-              context.instanceContext::failed);
+          isSameNode(upstreamEndPoint)
+              ? DATA_BLOCK_MANAGER.createLocalSourceHandle(
+                  localInstanceId.toThrift(),
+                  node.getPlanNodeId().getId(),
+                  remoteInstanceId.toThrift(),
+                  context.instanceContext::failed)
+              : DATA_BLOCK_MANAGER.createSourceHandle(
+                  localInstanceId.toThrift(),
+                  node.getPlanNodeId().getId(),
+                  upstreamEndPoint,
+                  remoteInstanceId.toThrift(),
+                  context.instanceContext::failed);
       return new ExchangeOperator(operatorContext, sourceHandle, 
node.getUpstreamPlanNodeId());
     }
 
@@ -1071,17 +1080,30 @@ public class LocalExecutionPlanner {
 
       FragmentInstanceId localInstanceId = context.instanceContext.getId();
       FragmentInstanceId targetInstanceId = node.getDownStreamInstanceId();
+
+      TEndPoint downStreamEndPoint = node.getDownStreamEndpoint();
+
       ISinkHandle sinkHandle =
-          DATA_BLOCK_MANAGER.createSinkHandle(
-              localInstanceId.toThrift(),
-              node.getDownStreamEndpoint(),
-              targetInstanceId.toThrift(),
-              node.getDownStreamPlanNodeId().getId(),
-              context.instanceContext);
+          isSameNode(downStreamEndPoint)
+              ? DATA_BLOCK_MANAGER.createLocalSinkHandle(
+                  localInstanceId.toThrift(),
+                  targetInstanceId.toThrift(),
+                  node.getDownStreamPlanNodeId().getId(),
+                  context.instanceContext)
+              : DATA_BLOCK_MANAGER.createSinkHandle(
+                  localInstanceId.toThrift(),
+                  downStreamEndPoint,
+                  targetInstanceId.toThrift(),
+                  node.getDownStreamPlanNodeId().getId(),
+                  context.instanceContext);
       context.setSinkHandle(sinkHandle);
       return child;
     }
 
+    private boolean isSameNode(TEndPoint endPoint) {
+      return endPoint.equals(LOCAL_HOST_DATA_BLOCK_ENDPOINT);
+    }
+
     @Override
     public Operator visitSchemaFetchMerge(
         SchemaFetchMergeNode node, LocalExecutionPlanContext context) {

Reply via email to