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),