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

xingtanzjr 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 149d88e7bf Improve auto create schema (#6295)
149d88e7bf is described below

commit 149d88e7bf3f47616b80e6c1d5a96d9653cffed9
Author: Marcos_Zyk <[email protected]>
AuthorDate: Thu Jun 16 10:56:53 2022 +0800

    Improve auto create schema (#6295)
---
 .../metadata/MeasurementAlreadyExistException.java |  42 ++++++
 .../iotdb/db/metadata/LocalSchemaProcessor.java    |   5 +-
 .../db/metadata/mtree/MTreeBelowSGMemoryImpl.java  |  19 ++-
 .../metadata/visitor/SchemaExecutionVisitor.java   | 163 ++++++++++++++++++---
 .../apache/iotdb/db/mpp/plan/analyze/Analyzer.java |  12 +-
 .../db/mpp/plan/analyze/ClusterSchemaFetcher.java  | 127 ++++++++--------
 .../iotdb/db/mpp/plan/constant/StatementType.java  |   3 +-
 .../iotdb/db/mpp/plan/planner/LogicalPlanner.java  |  28 ++--
 .../mpp/plan/planner/plan/node/PlanNodeType.java   |   6 +-
 .../db/mpp/plan/planner/plan/node/PlanVisitor.java |   5 +
 .../write/InternalCreateTimeSeriesNode.java        | 155 ++++++++++++++++++++
 .../db/mpp/plan/statement/StatementVisitor.java    |   8 +-
 .../InternalCreateTimeSeriesStatement.java}        |  40 ++++-
 .../java/org/apache/iotdb/rpc/TSStatusCode.java    |   1 +
 14 files changed, 491 insertions(+), 123 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/exception/metadata/MeasurementAlreadyExistException.java
 
b/server/src/main/java/org/apache/iotdb/db/exception/metadata/MeasurementAlreadyExistException.java
new file mode 100644
index 0000000000..3ca31b2deb
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/exception/metadata/MeasurementAlreadyExistException.java
@@ -0,0 +1,42 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ *
+ */
+
+package org.apache.iotdb.db.exception.metadata;
+
+import org.apache.iotdb.commons.exception.MetadataException;
+import org.apache.iotdb.db.metadata.path.MeasurementPath;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+public class MeasurementAlreadyExistException extends MetadataException {
+
+  private MeasurementPath measurementPath;
+
+  public MeasurementAlreadyExistException(String path, MeasurementPath 
measurementPath) {
+    super(
+        String.format("Path [%s] already exist", path),
+        TSStatusCode.MEASUREMENT_ALREADY_EXIST.getStatusCode());
+    this.isUserException = true;
+    this.measurementPath = measurementPath;
+  }
+
+  public MeasurementPath getMeasurementPath() {
+    return measurementPath;
+  }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/LocalSchemaProcessor.java 
b/server/src/main/java/org/apache/iotdb/db/metadata/LocalSchemaProcessor.java
index ee4f0be5f9..e365c90b57 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/LocalSchemaProcessor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/LocalSchemaProcessor.java
@@ -25,6 +25,7 @@ import org.apache.iotdb.commons.utils.TestOnly;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.metadata.AliasAlreadyExistException;
+import org.apache.iotdb.db.exception.metadata.MeasurementAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.PathAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.PathNotExistException;
 import org.apache.iotdb.db.exception.metadata.StorageGroupNotSetException;
@@ -308,7 +309,9 @@ public class LocalSchemaProcessor {
     try {
       createTimeseries(
           new CreateTimeSeriesPlan(path, dataType, encoding, compressor, 
props, null, null, null));
-    } catch (PathAlreadyExistException | AliasAlreadyExistException e) {
+    } catch (PathAlreadyExistException
+        | AliasAlreadyExistException
+        | MeasurementAlreadyExistException e) {
       if (logger.isDebugEnabled()) {
         logger.debug(
             "Ignore PathAlreadyExistException and AliasAlreadyExistException 
when Concurrent inserting"
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
index 7732cc5dbc..ab1f2d0e37 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/mtree/MTreeBelowSGMemoryImpl.java
@@ -26,6 +26,7 @@ import org.apache.iotdb.commons.utils.PathUtils;
 import org.apache.iotdb.db.exception.metadata.AliasAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.AlignedTimeseriesException;
 import org.apache.iotdb.db.exception.metadata.MNodeTypeMismatchException;
+import org.apache.iotdb.db.exception.metadata.MeasurementAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.PathAlreadyExistException;
 import org.apache.iotdb.db.exception.metadata.PathNotExistException;
 import 
org.apache.iotdb.db.exception.metadata.template.TemplateImcompatibeException;
@@ -174,7 +175,13 @@ public class MTreeBelowSGMemoryImpl implements 
IMTreeBelowSG {
       }
 
       if (device.hasChild(leafName)) {
-        throw new PathAlreadyExistException(path.getFullPath());
+        IMNode node = device.getChild(leafName);
+        if (node.isMeasurement()) {
+          throw new MeasurementAlreadyExistException(
+              path.getFullPath(), 
node.getAsMeasurementMNode().getMeasurementPath());
+        } else {
+          throw new PathAlreadyExistException(path.getFullPath());
+        }
       }
 
       if (upperTemplate != null
@@ -244,7 +251,15 @@ public class MTreeBelowSGMemoryImpl implements 
IMTreeBelowSG {
     synchronized (this) {
       for (int i = 0; i < measurements.size(); i++) {
         if (device.hasChild(measurements.get(i))) {
-          throw new PathAlreadyExistException(devicePath.getFullPath() + "." + 
measurements.get(i));
+          IMNode node = device.getChild(measurements.get(i));
+          if (node.isMeasurement()) {
+            throw new MeasurementAlreadyExistException(
+                devicePath.getFullPath() + "." + measurements.get(i),
+                node.getAsMeasurementMNode().getMeasurementPath());
+          } else {
+            throw new PathAlreadyExistException(
+                devicePath.getFullPath() + "." + measurements.get(i));
+          }
         }
         if (aliasList != null && aliasList.get(i) != null && 
device.hasChild(aliasList.get(i))) {
           throw new AliasAlreadyExistException(
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
index 3e5cc1068b..58ec1975d1 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
@@ -23,6 +23,8 @@ import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.conf.IoTDBConstant;
 import org.apache.iotdb.commons.exception.MetadataException;
 import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.db.exception.metadata.MeasurementAlreadyExistException;
+import org.apache.iotdb.db.metadata.path.MeasurementPath;
 import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor;
@@ -30,6 +32,7 @@ import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.AlterTimeSe
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateAlignedTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateMultiTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateTimeSeriesNode;
+import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InternalCreateTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.MeasurementGroup;
 import org.apache.iotdb.db.qp.physical.PhysicalPlan;
 import org.apache.iotdb.db.qp.physical.sys.CreateAlignedTimeSeriesPlan;
@@ -37,10 +40,15 @@ import 
org.apache.iotdb.db.qp.physical.sys.CreateTimeSeriesPlan;
 import org.apache.iotdb.rpc.RpcUtils;
 import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.tsfile.exception.NotImplementedException;
+import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.ByteArrayOutputStream;
+import java.io.DataOutputStream;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
@@ -89,26 +97,9 @@ public class SchemaExecutionVisitor extends 
PlanVisitor<TSStatus, ISchemaRegion>
       size = measurementGroup.getMeasurements().size();
       // todo implement batch creation of one device in SchemaRegion
       for (int i = 0; i < size; i++) {
-        CreateTimeSeriesPlan plan =
-            new CreateTimeSeriesPlan(
-                
devicePath.concatNode(measurementGroup.getMeasurements().get(i)),
-                measurementGroup.getDataTypes().get(i),
-                measurementGroup.getEncodings().get(i),
-                measurementGroup.getCompressors().get(i),
-                measurementGroup.getPropsList() == null
-                    ? null
-                    : measurementGroup.getPropsList().get(i),
-                measurementGroup.getTagsList() == null
-                    ? null
-                    : measurementGroup.getTagsList().get(i),
-                measurementGroup.getAttributesList() == null
-                    ? null
-                    : measurementGroup.getAttributesList().get(i),
-                measurementGroup.getAliasList() == null
-                    ? null
-                    : measurementGroup.getAliasList().get(i));
         try {
-          schemaRegion.createTimeseries(plan, -1);
+          schemaRegion.createTimeseries(
+              transformToCreateTimeSeriesPlan(devicePath, measurementGroup, 
i), -1);
         } catch (MetadataException e) {
           logger.error("{}: MetaData error: ", IoTDBConstant.GLOBAL_DB_NAME, 
e);
           failingStatus.add(RpcUtils.getStatus(e.getErrorCode(), 
e.getMessage()));
@@ -122,6 +113,140 @@ public class SchemaExecutionVisitor extends 
PlanVisitor<TSStatus, ISchemaRegion>
     return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS, "Execute 
successfully");
   }
 
+  private CreateTimeSeriesPlan transformToCreateTimeSeriesPlan(
+      PartialPath devicePath, MeasurementGroup measurementGroup, int index) {
+    return new CreateTimeSeriesPlan(
+        devicePath.concatNode(measurementGroup.getMeasurements().get(index)),
+        measurementGroup.getDataTypes().get(index),
+        measurementGroup.getEncodings().get(index),
+        measurementGroup.getCompressors().get(index),
+        measurementGroup.getPropsList() == null ? null : 
measurementGroup.getPropsList().get(index),
+        measurementGroup.getTagsList() == null ? null : 
measurementGroup.getTagsList().get(index),
+        measurementGroup.getAttributesList() == null
+            ? null
+            : measurementGroup.getAttributesList().get(index),
+        measurementGroup.getAliasList() == null
+            ? null
+            : measurementGroup.getAliasList().get(index));
+  }
+
+  @Override
+  public TSStatus visitInternalCreateTimeSeries(
+      InternalCreateTimeSeriesNode node, ISchemaRegion schemaRegion) {
+    PartialPath devicePath = node.getDevicePath();
+    MeasurementGroup measurementGroup = node.getMeasurementGroup();
+
+    List<TSStatus> alreadyExistingTimeseries = new ArrayList<>();
+    List<TSStatus> failingStatus = new ArrayList<>();
+
+    if (node.isAligned()) {
+      executeInternalCreateAlignedTimeseries(
+          devicePath, measurementGroup, schemaRegion, 
alreadyExistingTimeseries, failingStatus);
+    } else {
+      executeInternalCreateTimeseries(
+          devicePath, measurementGroup, schemaRegion, 
alreadyExistingTimeseries, failingStatus);
+    }
+
+    if (!failingStatus.isEmpty()) {
+      return RpcUtils.getStatus(failingStatus);
+    }
+
+    if (!alreadyExistingTimeseries.isEmpty()) {
+      return RpcUtils.getStatus(alreadyExistingTimeseries);
+    }
+
+    return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS, "Execute 
successfully");
+  }
+
+  private void executeInternalCreateTimeseries(
+      PartialPath devicePath,
+      MeasurementGroup measurementGroup,
+      ISchemaRegion schemaRegion,
+      List<TSStatus> alreadyExistingTimeseries,
+      List<TSStatus> failingStatus) {
+
+    int size = measurementGroup.getMeasurements().size();
+    // todo implement batch creation of one device in SchemaRegion
+    for (int i = 0; i < size; i++) {
+      try {
+        schemaRegion.createTimeseries(
+            transformToCreateTimeSeriesPlan(devicePath, measurementGroup, i), 
-1);
+      } catch (MeasurementAlreadyExistException e) {
+        logger.info("There's no need to internal create timeseries. {}", 
e.getMessage());
+        alreadyExistingTimeseries.add(
+            RpcUtils.getStatus(
+                e.getErrorCode(), 
transformExistingTimeseriesToString(e.getMeasurementPath())));
+      } catch (MetadataException e) {
+        logger.error("{}: MetaData error: ", IoTDBConstant.GLOBAL_DB_NAME, e);
+        failingStatus.add(RpcUtils.getStatus(e.getErrorCode(), 
e.getMessage()));
+      }
+    }
+  }
+
+  private String transformExistingTimeseriesToString(MeasurementPath 
measurementPath) {
+    ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
+    DataOutputStream dataOutputStream = new 
DataOutputStream(byteArrayOutputStream);
+    try {
+      measurementPath.serialize(dataOutputStream);
+    } catch (IOException ignored) {
+      // this exception won't happen.
+    }
+    return byteArrayOutputStream.toString();
+  }
+
+  private void executeInternalCreateAlignedTimeseries(
+      PartialPath devicePath,
+      MeasurementGroup measurementGroup,
+      ISchemaRegion schemaRegion,
+      List<TSStatus> alreadyExistingTimeseries,
+      List<TSStatus> failingStatus) {
+    List<String> measurementList = measurementGroup.getMeasurements();
+    List<TSDataType> dataTypeList = measurementGroup.getDataTypes();
+    List<TSEncoding> encodingList = measurementGroup.getEncodings();
+    List<CompressionType> compressionTypeList = 
measurementGroup.getCompressors();
+    CreateAlignedTimeSeriesPlan createAlignedTimeSeriesPlan =
+        new CreateAlignedTimeSeriesPlan(
+            devicePath,
+            measurementList,
+            dataTypeList,
+            encodingList,
+            compressionTypeList,
+            null,
+            null,
+            null);
+
+    boolean shouldRetry = true;
+    while (shouldRetry) {
+      try {
+        schemaRegion.createAlignedTimeSeries(createAlignedTimeSeriesPlan);
+        shouldRetry = false;
+      } catch (MeasurementAlreadyExistException e) {
+        // the existence check will be executed before truly creation
+        logger.info("There's no need to internal create timeseries. {}", 
e.getMessage());
+        MeasurementPath measurementPath = e.getMeasurementPath();
+        alreadyExistingTimeseries.add(
+            RpcUtils.getStatus(
+                e.getErrorCode(), 
transformExistingTimeseriesToString(measurementPath)));
+
+        // remove the existing timeseries from plan
+        int index = measurementList.indexOf(measurementPath.getMeasurement());
+        measurementList.remove(index);
+        dataTypeList.remove(index);
+        encodingList.remove(index);
+        compressionTypeList.remove(index);
+
+        if (measurementList.isEmpty()) {
+          shouldRetry = false;
+        }
+
+      } catch (MetadataException e) {
+        logger.error("{}: MetaData error: ", IoTDBConstant.GLOBAL_DB_NAME, e);
+        failingStatus.add(RpcUtils.getStatus(e.getErrorCode(), 
e.getMessage()));
+        shouldRetry = false;
+      }
+    }
+  }
+
   @Override
   public TSStatus visitAlterTimeSeries(AlterTimeSeriesNode node, ISchemaRegion 
schemaRegion) {
     try {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
index 9cf383cd8a..ae40bfcd6f 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analyzer.java
@@ -58,6 +58,7 @@ import 
org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowsStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertTabletStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.QueryStatement;
+import 
org.apache.iotdb.db.mpp.plan.statement.internal.InternalCreateTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.LastPointFetchStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.SchemaFetchStatement;
 import org.apache.iotdb.db.mpp.plan.statement.literal.Literal;
@@ -69,7 +70,6 @@ import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CountStorageGroupStatemen
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CountTimeSeriesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateAlignedTimeSeriesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateMultiTimeSeriesStatement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesByDeviceStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.ShowChildNodesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.ShowChildPathsStatement;
@@ -973,20 +973,20 @@ public class Analyzer {
     }
 
     @Override
-    public Analysis visitCreateTimeseriesByDevice(
-        CreateTimeSeriesByDeviceStatement createTimeSeriesByDeviceStatement,
+    public Analysis visitInternalCreateTimeseries(
+        InternalCreateTimeSeriesStatement internalCreateTimeSeriesStatement,
         MPPQueryContext context) {
       context.setQueryType(QueryType.WRITE);
 
       Analysis analysis = new Analysis();
-      analysis.setStatement(createTimeSeriesByDeviceStatement);
+      analysis.setStatement(internalCreateTimeSeriesStatement);
 
       SchemaPartition schemaPartitionInfo;
       schemaPartitionInfo =
           partitionFetcher.getOrCreateSchemaPartition(
               new PathPatternTree(
-                  createTimeSeriesByDeviceStatement.getDevicePath(),
-                  createTimeSeriesByDeviceStatement.getMeasurements()));
+                  internalCreateTimeSeriesStatement.getDevicePath(),
+                  internalCreateTimeSeriesStatement.getMeasurements()));
       analysis.setSchemaPartitionInfo(schemaPartitionInfo);
       return analysis;
     }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
index d5e41a50b7..c314051244 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/ClusterSchemaFetcher.java
@@ -26,15 +26,15 @@ import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.metadata.cache.DataNodeSchemaCache;
+import org.apache.iotdb.db.metadata.path.MeasurementPath;
+import org.apache.iotdb.db.metadata.path.PathDeserializeUtil;
 import org.apache.iotdb.db.mpp.common.schematree.DeviceSchemaInfo;
 import org.apache.iotdb.db.mpp.common.schematree.PathPatternTree;
 import org.apache.iotdb.db.mpp.common.schematree.SchemaTree;
 import org.apache.iotdb.db.mpp.plan.Coordinator;
 import org.apache.iotdb.db.mpp.plan.execution.ExecutionResult;
-import org.apache.iotdb.db.mpp.plan.statement.Statement;
+import 
org.apache.iotdb.db.mpp.plan.statement.internal.InternalCreateTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.SchemaFetchStatement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateAlignedTimeSeriesStatement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesByDeviceStatement;
 import org.apache.iotdb.db.query.control.SessionManager;
 import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
@@ -52,9 +52,12 @@ import io.airlift.concurrent.SetThreadName;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.List;
 import java.util.Map;
 import java.util.Optional;
+import java.util.Set;
+import java.util.stream.Collectors;
 
 import static 
org.apache.iotdb.db.utils.EncodingInferenceUtils.getDefaultEncoding;
 
@@ -245,34 +248,8 @@ public class ClusterSchemaFetcher implements 
ISchemaFetcher {
       return new SchemaTree();
     }
 
-    internalCreateTimeseries(
+    return internalCreateTimeseries(
         devicePath, missingMeasurements, dataTypesOfMissingMeasurement, 
isAligned);
-
-    SchemaTree reFetchSchemaTree =
-        fetchSchema(new PathPatternTree(devicePath, missingMeasurements));
-
-    Pair<List<String>, List<TSDataType>> recheckResult =
-        checkMissingMeasurements(
-            reFetchSchemaTree,
-            devicePath,
-            missingMeasurements.toArray(new String[0]),
-            dataTypesOfMissingMeasurement.toArray(new TSDataType[0]));
-
-    missingMeasurements = recheckResult.left;
-    if (!missingMeasurements.isEmpty()) {
-      StringBuilder stringBuilder = new StringBuilder();
-      stringBuilder.append("(");
-      for (String missingMeasurement : missingMeasurements) {
-        stringBuilder.append(missingMeasurement).append(" ");
-      }
-      stringBuilder.append(")");
-      throw new RuntimeException(
-          String.format(
-              "Failed to auto create schema, devicePath: %s, measurements: %s",
-              devicePath.getFullPath(), stringBuilder));
-    }
-
-    return reFetchSchemaTree;
   }
 
   private Pair<List<String>, List<TSDataType>> checkMissingMeasurements(
@@ -299,65 +276,81 @@ public class ClusterSchemaFetcher implements 
ISchemaFetcher {
     return new Pair<>(missingMeasurements, dataTypesOfMissingMeasurement);
   }
 
-  private void internalCreateTimeseries(
+  private SchemaTree internalCreateTimeseries(
       PartialPath devicePath,
       List<String> measurements,
       List<TSDataType> tsDataTypes,
       boolean isAligned) {
 
-    if (isAligned) {
-      CreateAlignedTimeSeriesStatement createAlignedTimeSeriesStatement =
-          new CreateAlignedTimeSeriesStatement();
-      createAlignedTimeSeriesStatement.setDevicePath(devicePath);
-      createAlignedTimeSeriesStatement.setMeasurements(measurements);
-      createAlignedTimeSeriesStatement.setDataTypes(tsDataTypes);
-      List<TSEncoding> encodings = new ArrayList<>();
-      List<CompressionType> compressors = new ArrayList<>();
-      for (TSDataType dataType : tsDataTypes) {
-        encodings.add(getDefaultEncoding(dataType));
-        
compressors.add(TSFileDescriptor.getInstance().getConfig().getCompressor());
-      }
-      createAlignedTimeSeriesStatement.setEncodings(encodings);
-      createAlignedTimeSeriesStatement.setCompressors(compressors);
-      createAlignedTimeSeriesStatement.setAliasList(null);
+    List<TSEncoding> encodings = new ArrayList<>();
+    List<CompressionType> compressors = new ArrayList<>();
+    for (TSDataType dataType : tsDataTypes) {
+      encodings.add(getDefaultEncoding(dataType));
+      
compressors.add(TSFileDescriptor.getInstance().getConfig().getCompressor());
+    }
 
-      executeCreateStatement(createAlignedTimeSeriesStatement);
-    } else {
+    List<MeasurementPath> measurementPathList =
+        executeInternalCreateTimeseriesStatement(
+            new InternalCreateTimeSeriesStatement(
+                devicePath, measurements, tsDataTypes, encodings, compressors, 
isAligned));
 
-      executeCreateTimeseriesByDeviceStatement(
-          new CreateTimeSeriesByDeviceStatement(devicePath, measurements, 
tsDataTypes));
-    }
-  }
+    Set<Integer> alreadyExistingMeasurementIndexSet =
+        measurementPathList.stream()
+            .map(o -> measurements.indexOf(o.getMeasurement()))
+            .collect(Collectors.toSet());
 
-  private void executeCreateStatement(Statement statement) {
-    long queryId = SessionManager.getInstance().requestQueryId(false);
-    ExecutionResult executionResult =
-        coordinator.execute(statement, queryId, null, "", partitionFetcher, 
this);
-    // TODO: throw exception
-    int statusCode = executionResult.status.getCode();
-    if (statusCode != TSStatusCode.SUCCESS_STATUS.getStatusCode()
-        && statusCode != 
TSStatusCode.PATH_ALREADY_EXIST_ERROR.getStatusCode()) {
-      throw new RuntimeException("cannot auto create schema, status is: " + 
executionResult.status);
+    SchemaTree schemaTree = new SchemaTree();
+    schemaTree.appendMeasurementPaths(measurementPathList);
+
+    for (int i = 0, size = measurements.size(); i < size; i++) {
+      if (alreadyExistingMeasurementIndexSet.contains(i)) {
+        continue;
+      }
+
+      schemaTree.appendSingleMeasurement(
+          devicePath.concatNode(measurements.get(i)),
+          new MeasurementSchema(
+              measurements.get(i), tsDataTypes.get(i), encodings.get(i), 
compressors.get(i)),
+          null,
+          isAligned);
     }
+
+    return schemaTree;
   }
 
-  private void executeCreateTimeseriesByDeviceStatement(
-      CreateTimeSeriesByDeviceStatement statement) {
+  private List<MeasurementPath> executeInternalCreateTimeseriesStatement(
+      InternalCreateTimeSeriesStatement statement) {
     long queryId = SessionManager.getInstance().requestQueryId(false);
     ExecutionResult executionResult =
         coordinator.execute(statement, queryId, null, "", partitionFetcher, 
this);
     // TODO: throw exception
     int statusCode = executionResult.status.getCode();
     if (statusCode == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      return;
+      return Collections.emptyList();
     }
 
+    List<String> failedCreationList = new ArrayList<>();
+    List<MeasurementPath> alreadyExistingMeasurements = new ArrayList<>();
     for (TSStatus subStatus : executionResult.status.subStatus) {
-      if (subStatus.code != 
TSStatusCode.PATH_ALREADY_EXIST_ERROR.getStatusCode()) {
-        throw new RuntimeException(
-            "cannot auto create schema, status is: " + executionResult.status);
+      if (subStatus.code == 
TSStatusCode.MEASUREMENT_ALREADY_EXIST.getStatusCode()) {
+        alreadyExistingMeasurements.add(
+            (MeasurementPath)
+                PathDeserializeUtil.deserialize(
+                    ByteBuffer.wrap(subStatus.getMessage().getBytes())));
+      } else {
+        failedCreationList.add(subStatus.message);
       }
     }
+
+    if (!failedCreationList.isEmpty()) {
+      StringBuilder stringBuilder = new StringBuilder();
+      for (String message : failedCreationList) {
+        stringBuilder.append(message).append("\n");
+      }
+      throw new RuntimeException(String.format("Failed to auto create schema\n 
%s", stringBuilder));
+    }
+
+    return alreadyExistingMeasurements;
   }
 
   @Override
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/constant/StatementType.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/constant/StatementType.java
index 42d1cd39a7..c4e14cc331 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/constant/StatementType.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/constant/StatementType.java
@@ -53,7 +53,6 @@ public enum StatementType {
   DELETE_STORAGE_GROUP,
   CREATE_TIMESERIES,
   CREATE_ALIGNED_TIMESERIES,
-  CREATE_TIMESERIES_BY_DEVICE,
   CREATE_MULTI_TIMESERIES,
   DELETE_TIMESERIES,
   ALTER_TIMESERIES,
@@ -134,7 +133,7 @@ public enum StatementType {
   SHOW_QUERY_RESOURCE,
 
   FETCH_SCHEMA,
-  FETCH_SCHEMA_WITH_AUTO_CREATE,
+  INTERNAL_CREATE_TIMESERIES,
 
   COUNT
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
index 0364111bb6..22bc078f3b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanner.java
@@ -29,6 +29,7 @@ import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.AlterTimeSe
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateAlignedTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateMultiTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateTimeSeriesNode;
+import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InternalCreateTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.MeasurementGroup;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.DeleteDataNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertMultiTabletsNode;
@@ -46,6 +47,7 @@ import 
org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowsOfOneDeviceStatemen
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowsStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertTabletStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.QueryStatement;
+import 
org.apache.iotdb.db.mpp.plan.statement.internal.InternalCreateTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.LastPointFetchStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.SchemaFetchStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.AlterTimeSeriesStatement;
@@ -55,24 +57,19 @@ import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CountNodesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CountTimeSeriesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateAlignedTimeSeriesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateMultiTimeSeriesStatement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesByDeviceStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.ShowChildNodesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.ShowChildPathsStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.ShowDevicesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.ShowTimeSeriesStatement;
-import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
 
 import java.util.ArrayList;
-import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.stream.Collectors;
 
-import static 
org.apache.iotdb.db.utils.EncodingInferenceUtils.getDefaultEncoding;
-
 /** Generate a logical plan for the statement. */
 public class LogicalPlanner {
 
@@ -362,24 +359,25 @@ public class LogicalPlanner {
     }
 
     @Override
-    public PlanNode visitCreateTimeseriesByDevice(
-        CreateTimeSeriesByDeviceStatement createTimeSeriesByDeviceStatement,
+    public PlanNode visitInternalCreateTimeseries(
+        InternalCreateTimeSeriesStatement internalCreateTimeSeriesStatement,
         MPPQueryContext context) {
-      int size = createTimeSeriesByDeviceStatement.getMeasurements().size();
+      int size = internalCreateTimeSeriesStatement.getMeasurements().size();
 
       MeasurementGroup measurementGroup = new MeasurementGroup();
       for (int i = 0; i < size; i++) {
         measurementGroup.addMeasurement(
-            createTimeSeriesByDeviceStatement.getMeasurements().get(i),
-            createTimeSeriesByDeviceStatement.getTsDataTypes().get(i),
-            
getDefaultEncoding(createTimeSeriesByDeviceStatement.getTsDataTypes().get(i)),
-            TSFileDescriptor.getInstance().getConfig().getCompressor());
+            internalCreateTimeSeriesStatement.getMeasurements().get(i),
+            internalCreateTimeSeriesStatement.getTsDataTypes().get(i),
+            internalCreateTimeSeriesStatement.getEncodings().get(i),
+            internalCreateTimeSeriesStatement.getCompressors().get(i));
       }
 
-      return new CreateMultiTimeSeriesNode(
+      return new InternalCreateTimeSeriesNode(
           context.getQueryId().genPlanNodeId(),
-          Collections.singletonMap(
-              createTimeSeriesByDeviceStatement.getDevicePath(), 
measurementGroup));
+          internalCreateTimeSeriesStatement.getDevicePath(),
+          measurementGroup,
+          internalCreateTimeSeriesStatement.isAligned());
     }
 
     @Override
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
index d95329a3e3..a7e50182a4 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
@@ -38,6 +38,7 @@ import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateAlign
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateMultiTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.DeleteTimeSeriesNode;
+import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InternalCreateTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InvalidateSchemaCacheNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.AggregationNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceMergeNode;
@@ -124,7 +125,8 @@ public enum PlanNodeType {
   LAST_QUERY_SCAN((short) 46),
   ALIGNED_LAST_QUERY_SCAN((short) 47),
   LAST_QUERY_MERGE((short) 48),
-  NODE_PATHS_COUNT((short) 49);
+  NODE_PATHS_COUNT((short) 49),
+  INTERNAL_CREATE_TIMESERIES((short) 50);
 
   private final short nodeType;
 
@@ -254,6 +256,8 @@ public enum PlanNodeType {
         return LastQueryMergeNode.deserialize(buffer);
       case 49:
         return NodePathsCountNode.deserialize(buffer);
+      case 50:
+        return InternalCreateTimeSeriesNode.deserialize(buffer);
       default:
         throw new IllegalArgumentException("Invalid node type: " + nodeType);
     }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
index 068e753aa7..f1f9092b86 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
@@ -38,6 +38,7 @@ import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateAlign
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateMultiTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.CreateTimeSeriesNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.DeleteTimeSeriesNode;
+import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InternalCreateTimeSeriesNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.AggregationNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceMergeNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.process.DeviceViewNode;
@@ -271,4 +272,8 @@ public abstract class PlanVisitor<R, C> {
   public R visitDeleteData(DeleteDataNode node, C context) {
     return visitPlan(node, context);
   }
+
+  public R visitInternalCreateTimeSeries(InternalCreateTimeSeriesNode node, C 
context) {
+    return visitPlan(node, context);
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/write/InternalCreateTimeSeriesNode.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/write/InternalCreateTimeSeriesNode.java
new file mode 100644
index 0000000000..ee1e5bb33d
--- /dev/null
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/write/InternalCreateTimeSeriesNode.java
@@ -0,0 +1,155 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write;
+
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.db.metadata.path.PathDeserializeUtil;
+import org.apache.iotdb.db.mpp.plan.analyze.Analysis;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeType;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.WritePlanNode;
+import org.apache.iotdb.tsfile.exception.NotImplementedException;
+import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
+
+import com.google.common.collect.ImmutableList;
+
+import java.io.DataOutputStream;
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+
+public class InternalCreateTimeSeriesNode extends WritePlanNode {
+
+  private PartialPath devicePath;
+  private MeasurementGroup measurementGroup;
+  private boolean isAligned;
+
+  private TRegionReplicaSet regionReplicaSet;
+
+  public InternalCreateTimeSeriesNode(
+      PlanNodeId id, PartialPath devicePath, MeasurementGroup 
measurementGroup, boolean isAligned) {
+    super(id);
+    this.devicePath = devicePath;
+    this.measurementGroup = measurementGroup;
+    this.isAligned = isAligned;
+  }
+
+  public PartialPath getDevicePath() {
+    return devicePath;
+  }
+
+  public MeasurementGroup getMeasurementGroup() {
+    return measurementGroup;
+  }
+
+  public boolean isAligned() {
+    return isAligned;
+  }
+
+  @Override
+  public TRegionReplicaSet getRegionReplicaSet() {
+    return regionReplicaSet;
+  }
+
+  public void setRegionReplicaSet(TRegionReplicaSet regionReplicaSet) {
+    this.regionReplicaSet = regionReplicaSet;
+  }
+
+  @Override
+  public List<PlanNode> getChildren() {
+    return new ArrayList<>();
+  }
+
+  @Override
+  public void addChild(PlanNode child) {}
+
+  @Override
+  public PlanNode clone() {
+    throw new NotImplementedException("Clone of InternalCreateTimeSeriesNode 
is not implemented");
+  }
+
+  @Override
+  public int allowedChildCount() {
+    return NO_CHILD_ALLOWED;
+  }
+
+  @Override
+  public List<String> getOutputColumnNames() {
+    return null;
+  }
+
+  @Override
+  public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
+    return visitor.visitInternalCreateTimeSeries(this, context);
+  }
+
+  @Override
+  protected void serializeAttributes(ByteBuffer byteBuffer) {
+    PlanNodeType.INTERNAL_CREATE_TIMESERIES.serialize(byteBuffer);
+    devicePath.serialize(byteBuffer);
+    measurementGroup.serialize(byteBuffer);
+    ReadWriteIOUtils.write(isAligned, byteBuffer);
+  }
+
+  @Override
+  protected void serializeAttributes(DataOutputStream stream) throws 
IOException {
+    PlanNodeType.INTERNAL_CREATE_TIMESERIES.serialize(stream);
+    devicePath.serialize(stream);
+    measurementGroup.serialize(stream);
+    ReadWriteIOUtils.write(isAligned, stream);
+  }
+
+  public static InternalCreateTimeSeriesNode deserialize(ByteBuffer 
byteBuffer) {
+    PartialPath devicePath = (PartialPath) 
PathDeserializeUtil.deserialize(byteBuffer);
+    MeasurementGroup measurementGroup = new MeasurementGroup();
+    measurementGroup.deserialize(byteBuffer);
+    boolean isAligned = ReadWriteIOUtils.readBool(byteBuffer);
+    PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer);
+    return new InternalCreateTimeSeriesNode(planNodeId, devicePath, 
measurementGroup, isAligned);
+  }
+
+  @Override
+  public List<WritePlanNode> splitByPartition(Analysis analysis) {
+    TRegionReplicaSet regionReplicaSet =
+        
analysis.getSchemaPartitionInfo().getSchemaRegionReplicaSet(devicePath.getFullPath());
+    setRegionReplicaSet(regionReplicaSet);
+    return ImmutableList.of(this);
+  }
+
+  @Override
+  public boolean equals(Object o) {
+    if (this == o) return true;
+    if (o == null || getClass() != o.getClass()) return false;
+    if (!super.equals(o)) return false;
+    InternalCreateTimeSeriesNode that = (InternalCreateTimeSeriesNode) o;
+    return Objects.equals(devicePath, that.devicePath)
+        && Objects.equals(measurementGroup, that.measurementGroup);
+  }
+
+  @Override
+  public int hashCode() {
+    return Objects.hash(super.hashCode(), devicePath, measurementGroup);
+  }
+}
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/StatementVisitor.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/StatementVisitor.java
index b878a7038a..601ca5e260 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/StatementVisitor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/StatementVisitor.java
@@ -27,6 +27,7 @@ import 
org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowsStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.InsertTabletStatement;
 import org.apache.iotdb.db.mpp.plan.statement.crud.QueryStatement;
+import 
org.apache.iotdb.db.mpp.plan.statement.internal.InternalCreateTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.LastPointFetchStatement;
 import org.apache.iotdb.db.mpp.plan.statement.internal.SchemaFetchStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.AlterTimeSeriesStatement;
@@ -38,7 +39,6 @@ import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CountTimeSeriesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateAlignedTimeSeriesStatement;
 import org.apache.iotdb.db.mpp.plan.statement.metadata.CreateFunctionStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateMultiTimeSeriesStatement;
-import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesByDeviceStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateTimeSeriesStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.DeleteStorageGroupStatement;
 import 
org.apache.iotdb.db.mpp.plan.statement.metadata.DeleteTimeSeriesStatement;
@@ -93,9 +93,9 @@ public abstract class StatementVisitor<R, C> {
   }
 
   // Create Timeseries by device
-  public R visitCreateTimeseriesByDevice(
-      CreateTimeSeriesByDeviceStatement createTimeSeriesByDeviceStatement, C 
context) {
-    return visitStatement(createTimeSeriesByDeviceStatement, context);
+  public R visitInternalCreateTimeseries(
+      InternalCreateTimeSeriesStatement internalCreateTimeSeriesStatement, C 
context) {
+    return visitStatement(internalCreateTimeSeriesStatement, context);
   }
 
   // Create Multi Timeseries
diff --git 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/CreateTimeSeriesByDeviceStatement.java
 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/internal/InternalCreateTimeSeriesStatement.java
similarity index 63%
rename from 
server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/CreateTimeSeriesByDeviceStatement.java
rename to 
server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/internal/InternalCreateTimeSeriesStatement.java
index 1bb8f08679..1db7154d76 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/CreateTimeSeriesByDeviceStatement.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/internal/InternalCreateTimeSeriesStatement.java
@@ -17,31 +17,47 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.mpp.plan.statement.metadata;
+package org.apache.iotdb.db.mpp.plan.statement.internal;
 
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.db.mpp.plan.constant.StatementType;
 import org.apache.iotdb.db.mpp.plan.statement.Statement;
 import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor;
+import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
 
+import java.util.ArrayList;
 import java.util.List;
 import java.util.stream.Collectors;
 
 // This is only used for auto creation while inserting data
-public class CreateTimeSeriesByDeviceStatement extends Statement {
+public class InternalCreateTimeSeriesStatement extends Statement {
 
   private PartialPath devicePath;
   private List<String> measurements;
+
   private List<TSDataType> tsDataTypes;
+  private List<TSEncoding> encodings = new ArrayList<>();
+  private List<CompressionType> compressors = new ArrayList<>();
+
+  private boolean isAligned;
 
-  public CreateTimeSeriesByDeviceStatement(
-      PartialPath devicePath, List<String> measurements, List<TSDataType> 
tsDataTypes) {
+  public InternalCreateTimeSeriesStatement(
+      PartialPath devicePath,
+      List<String> measurements,
+      List<TSDataType> tsDataTypes,
+      List<TSEncoding> encodings,
+      List<CompressionType> compressors,
+      boolean isAligned) {
     super();
-    setType(StatementType.CREATE_TIMESERIES_BY_DEVICE);
+    setType(StatementType.INTERNAL_CREATE_TIMESERIES);
     this.devicePath = devicePath;
     this.measurements = measurements;
     this.tsDataTypes = tsDataTypes;
+    this.encodings = encodings;
+    this.compressors = compressors;
+    this.isAligned = isAligned;
   }
 
   public PartialPath getDevicePath() {
@@ -56,6 +72,18 @@ public class CreateTimeSeriesByDeviceStatement extends 
Statement {
     return tsDataTypes;
   }
 
+  public List<TSEncoding> getEncodings() {
+    return encodings;
+  }
+
+  public List<CompressionType> getCompressors() {
+    return compressors;
+  }
+
+  public boolean isAligned() {
+    return isAligned;
+  }
+
   @Override
   public List<? extends PartialPath> getPaths() {
     return measurements.stream().map(o -> 
devicePath.concatNode(o)).collect(Collectors.toList());
@@ -63,6 +91,6 @@ public class CreateTimeSeriesByDeviceStatement extends 
Statement {
 
   @Override
   public <R, C> R accept(StatementVisitor<R, C> visitor, C context) {
-    return visitor.visitCreateTimeseriesByDevice(this, context);
+    return visitor.visitInternalCreateTimeseries(this, context);
   }
 }
diff --git a/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java 
b/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
index 90a36ede1d..f0bd4f94d7 100644
--- a/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
+++ b/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
@@ -67,6 +67,7 @@ public enum TSStatusCode {
   PIPE_ERROR(335),
   PIPESERVER_ERROR(336),
   SERIES_OVERFLOW(337),
+  MEASUREMENT_ALREADY_EXIST(338),
 
   EXECUTE_STATEMENT_ERROR(400),
   SQL_PARSE_ERROR(401),

Reply via email to