This is an automated email from the ASF dual-hosted git repository.
xingtanzjr pushed a commit to branch ml_0729_test
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/ml_0729_test by this push:
new 9c317d6546 change write to fake select host
9c317d6546 is described below
commit 9c317d65460f361b863d956e67e88395916d9082
Author: Jinrui.Zhang <[email protected]>
AuthorDate: Fri Jul 29 15:24:52 2022 +0800
change write to fake select host
---
.../distribution/SimpleFragmentParallelPlanner.java | 13 +------------
.../planner/distribution/WriteFragmentParallelPlanner.java | 14 ++++++++++++++
2 files changed, 15 insertions(+), 12 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
index f9398fc598..854bb15467 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/SimpleFragmentParallelPlanner.java
@@ -103,24 +103,13 @@ public class SimpleFragmentParallelPlanner implements
IFragmentParallelPlaner {
// redirected
// to another host when scheduling
fragmentInstance.setDataRegionAndHost(regionReplicaSet);
- fragmentInstance.setHostDataNode(fakeSelectDataNode(regionReplicaSet));
+ fragmentInstance.setHostDataNode(selectTargetDataNode(regionReplicaSet));
fragmentInstance.getFragment().setTypeProvider(analysis.getTypeProvider());
instanceMap.putIfAbsent(fragment.getId(), fragmentInstance);
fragmentInstanceList.add(fragmentInstance);
}
- private TDataNodeLocation fakeSelectDataNode(TRegionReplicaSet
regionReplicaSet) {
- String[] candidate = new String[] {"172.20.31.41", "172.20.31.42",
"172.20.31.43"};
- int targetIndex = regionReplicaSet.regionId.id % 3;
- for (TDataNodeLocation location : regionReplicaSet.getDataNodeLocations())
{
- if (location.internalEndPoint.getIp().equals(candidate[targetIndex])) {
- return location;
- }
- }
- return null;
- }
-
private TDataNodeLocation selectTargetDataNode(TRegionReplicaSet
regionReplicaSet) {
if (regionReplicaSet == null
|| regionReplicaSet.getDataNodeLocations() == null
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/WriteFragmentParallelPlanner.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/WriteFragmentParallelPlanner.java
index 4f043a05b9..cfc17b45cb 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/WriteFragmentParallelPlanner.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/distribution/WriteFragmentParallelPlanner.java
@@ -19,6 +19,8 @@
package org.apache.iotdb.db.mpp.plan.planner.distribution;
+import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
import org.apache.iotdb.db.mpp.common.MPPQueryContext;
import org.apache.iotdb.db.mpp.plan.analyze.Analysis;
import org.apache.iotdb.db.mpp.plan.planner.IFragmentParallelPlaner;
@@ -64,8 +66,20 @@ public class WriteFragmentParallelPlanner implements
IFragmentParallelPlaner {
queryContext.getQueryType(),
queryContext.getTimeOut());
instance.setDataRegionAndHost(split.getRegionReplicaSet());
+
instance.setHostDataNode(fakeSelectDataNode(split.getRegionReplicaSet()));
ret.add(instance);
}
return ret;
}
+
+ private TDataNodeLocation fakeSelectDataNode(TRegionReplicaSet
regionReplicaSet) {
+ String[] candidate = new String[] {"172.20.31.41", "172.20.31.42",
"172.20.31.43"};
+ int targetIndex = regionReplicaSet.regionId.id % 3;
+ for (TDataNodeLocation location : regionReplicaSet.getDataNodeLocations())
{
+ if (location.internalEndPoint.getIp().equals(candidate[targetIndex])) {
+ return location;
+ }
+ }
+ return null;
+ }
}