This is an automated email from the ASF dual-hosted git repository.
jackietien 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 20c679877d Fix the issue that sometimes the schema partition cannot be
calculated correctly (#6276)
20c679877d is described below
commit 20c679877dd5eac91753b053f658e438c993bbed
Author: Zhang.Jinrui <[email protected]>
AuthorDate: Wed Jun 15 21:09:49 2022 +0800
Fix the issue that sometimes the schema partition cannot be calculated
correctly (#6276)
---
.../iotdb/confignode/manager/ConfigManager.java | 94 ++++++++++++++--------
.../db/mpp/common/schematree/PathPatternTree.java | 6 +-
.../node/metedata/read/DevicesSchemaScanNode.java | 5 ++
3 files changed, 70 insertions(+), 35 deletions(-)
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
index e81c72d76c..75a0f1f04f 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
@@ -22,6 +22,9 @@ package org.apache.iotdb.confignode.manager;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.common.rpc.thrift.TSeriesPartitionSlot;
import org.apache.iotdb.commons.conf.CommonDescriptor;
+import org.apache.iotdb.commons.conf.IoTDBConstant;
+import org.apache.iotdb.commons.exception.IllegalPathException;
+import org.apache.iotdb.commons.partition.SchemaPartition;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.utils.AuthUtils;
import org.apache.iotdb.commons.utils.PathUtils;
@@ -65,6 +68,7 @@ import
org.apache.iotdb.confignode.rpc.thrift.TConfigNodeRegisterResp;
import org.apache.iotdb.confignode.rpc.thrift.TPermissionInfoResp;
import org.apache.iotdb.confignode.rpc.thrift.TStorageGroupSchema;
import org.apache.iotdb.consensus.common.DataSet;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.mpp.common.schematree.PathPatternTree;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -73,11 +77,15 @@ import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.stream.Collectors;
+
+import static
org.apache.iotdb.commons.conf.IoTDBConstant.MULTI_LEVEL_PATH_WILDCARD;
/** Entry of all management, AssignPartitionManager,AssignRegionManager. */
public class ConfigManager implements Manager {
@@ -267,56 +275,76 @@ public class ConfigManager implements Manager {
}
}
+ private List<TSeriesPartitionSlot> calculateRelatedSlot(
+ PartialPath path, PartialPath storageGroup) {
+ // The path contains `**`
+ if (path.getFullPath().contains(IoTDBConstant.MULTI_LEVEL_PATH_WILDCARD)) {
+ return new ArrayList<>();
+ }
+ // path doesn't contain * so the size of innerPathList should be 1
+ PartialPath innerPath = path.alterPrefixPath(storageGroup).get(0);
+ // The innerPath contains `*` and the only `*` is not in last level
+ if (innerPath.getDevice().contains(IoTDBConstant.ONE_LEVEL_PATH_WILDCARD))
{
+ return new ArrayList<>();
+ }
+ return Collections.singletonList(
+ getPartitionManager().getSeriesPartitionSlot(innerPath.getDevice()));
+ }
+
@Override
public DataSet getSchemaPartition(PathPatternTree patternTree) {
TSStatus status = confirmLeader();
if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- List<String> devicePaths = patternTree.findAllDevicePaths();
- List<String> storageGroups =
getClusterSchemaManager().getStorageGroupNames();
GetSchemaPartitionReq getSchemaPartitionReq = new
GetSchemaPartitionReq();
- Map<String, List<TSeriesPartitionSlot>> partitionSlotsMap = new
HashMap<>();
-
- boolean getAll = false;
- Set<String> getAllSet = new HashSet<>();
- for (String devicePath : devicePaths) {
- boolean matchStorageGroup = false;
- for (String storageGroup : storageGroups) {
- if (PathUtils.isStartWith(devicePath, storageGroup)) {
- matchStorageGroup = true;
- if (devicePath.contains("*")) {
- // Get all SchemaPartitions of this StorageGroup if the
devicePath contains "*"
- getAllSet.add(storageGroup);
- } else {
- // Get the specific SchemaPartition
- partitionSlotsMap
- .computeIfAbsent(storageGroup, key -> new ArrayList<>())
-
.add(getPartitionManager().getSeriesPartitionSlot(devicePath));
+ Map<String, Set<TSeriesPartitionSlot>> partitionSlotsMap = new
HashMap<>();
+ List<PartialPath> relatedPaths = patternTree.splitToPathList();
+ List<String> allStorageGroups =
getClusterSchemaManager().getStorageGroupNames();
+ Map<String, Boolean> scanAllRegions = new HashMap<>();
+ for (PartialPath path : relatedPaths) {
+ for (String storageGroup : allStorageGroups) {
+ try {
+ PartialPath storageGroupPath = new PartialPath(storageGroup);
+ if
(path.overlapWith(storageGroupPath.concatNode(MULTI_LEVEL_PATH_WILDCARD))
+ && !scanAllRegions.containsKey(storageGroup)) {
+ List<TSeriesPartitionSlot> relatedSlot =
calculateRelatedSlot(path, storageGroupPath);
+ if (relatedSlot.isEmpty()) {
+ scanAllRegions.put(storageGroup, true);
+ partitionSlotsMap.put(storageGroup, new HashSet<>());
+ } else {
+ partitionSlotsMap
+ .computeIfAbsent(storageGroup, k -> new HashSet<>())
+ .addAll(relatedSlot);
+ }
}
- break;
+ } catch (IllegalPathException e) {
+ // this line won't be reached in general
+ throw new RuntimeException(e);
}
}
- if (!matchStorageGroup && devicePath.contains("**")) {
- // Get all SchemaPartitions if there exists one devicePath that
contains "**"
- getAll = true;
- }
}
- if (getAll) {
- partitionSlotsMap = new HashMap<>();
- } else {
- for (String storageGroup : getAllSet) {
- partitionSlotsMap.put(storageGroup, new ArrayList<>());
- }
+ // return empty partition
+ if (partitionSlotsMap.isEmpty()) {
+ SchemaPartitionResp resp = new SchemaPartitionResp();
+ resp.setStatus(new
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
+ // TODO: (xingtanzjr) replace with new data structure
+ resp.setSchemaPartition(
+ new SchemaPartition(
+ new HashMap<>(),
+
IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionExecutorClass(),
+
IoTDBDescriptor.getInstance().getConfig().getSeriesPartitionSlotNum()));
}
- getSchemaPartitionReq.setPartitionSlotsMap(partitionSlotsMap);
+ getSchemaPartitionReq.setPartitionSlotsMap(
+ partitionSlotsMap.entrySet().stream()
+ .collect(Collectors.toMap(Map.Entry::getKey, e -> new
ArrayList<>(e.getValue()))));
SchemaPartitionResp resp =
(SchemaPartitionResp)
partitionManager.getSchemaPartition(getSchemaPartitionReq);
// TODO: Delete or hide this LOGGER before officially release.
LOGGER.info(
- "GetSchemaPartition interface receive devicePaths: {}, return
SchemaPartition: {}",
- devicePaths,
+ "GetSchemaPartition interface receive paths: {}, return
SchemaPartition: {}",
+ relatedPaths,
resp.getSchemaPartition().getSchemaPartitionMap());
return resp;
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/common/schematree/PathPatternTree.java
b/server/src/main/java/org/apache/iotdb/db/mpp/common/schematree/PathPatternTree.java
index 84047dd825..840eb3ccb6 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/common/schematree/PathPatternTree.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/common/schematree/PathPatternTree.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.mpp.common.schematree;
+import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.utils.TestOnly;
@@ -109,8 +110,9 @@ public class PathPatternTree {
PathPatternNode curNode, List<String> nodes, List<String>
pathPatternList) {
nodes.add(curNode.getName());
if (curNode.isLeaf()) {
- if (!curNode.getName().equals("**")) {
- pathPatternList.add(parseNodesToString(nodes.subList(0, nodes.size() -
1)));
+ if (!curNode.getName().equals(IoTDBConstant.MULTI_LEVEL_PATH_WILDCARD)) {
+ pathPatternList.add(
+ nodes.size() == 1 ? "" : parseNodesToString(nodes.subList(0,
nodes.size() - 1)));
} else {
pathPatternList.add(parseNodesToString(nodes));
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/DevicesSchemaScanNode.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/DevicesSchemaScanNode.java
index 9961d5c668..7672e61296 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/DevicesSchemaScanNode.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/read/DevicesSchemaScanNode.java
@@ -119,4 +119,9 @@ public class DevicesSchemaScanNode extends
SchemaQueryScanNode {
public int hashCode() {
return Objects.hash(super.hashCode(), hasSgCol);
}
+
+ @Override
+ public String toString() {
+ return String.format("DevicesSchemaScanNode-%s[Path: %s]",
getPlanNodeId(), path);
+ }
}