This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch forward_schema_validate in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 5090928fd4bb1e172d4e26ed8b615c0c0b010f7f Author: HTHou <[email protected]> AuthorDate: Fri May 5 16:04:28 2023 +0800 forward schema validate --- .../execution/executor/RegionWriteExecutor.java | 40 +--- .../apache/iotdb/db/mpp/plan/analyze/Analysis.java | 16 +- .../iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java | 54 +++-- .../mpp/plan/analyze/schema/SchemaValidator.java | 20 +- .../db/mpp/plan/execution/QueryExecution.java | 7 +- .../db/mpp/plan/planner/LogicalPlanVisitor.java | 5 + .../planner/plan/node/write/BatchInsertNode.java | 33 --- .../plan/node/write/InsertMultiTabletsNode.java | 18 +- .../plan/planner/plan/node/write/InsertNode.java | 131 ++---------- .../planner/plan/node/write/InsertRowNode.java | 150 ++------------ .../planner/plan/node/write/InsertRowsNode.java | 28 +-- .../plan/node/write/InsertRowsOfOneDeviceNode.java | 29 +-- .../planner/plan/node/write/InsertTabletNode.java | 119 +---------- .../scheduler/FragmentInstanceDispatcherImpl.java | 11 +- .../plan/statement/crud/InsertBaseStatement.java | 226 ++++++++++++++++++++- .../crud/InsertMultiTabletsStatement.java | 30 +++ .../plan/statement/crud/InsertRowStatement.java | 182 ++++++++++++++++- .../crud/InsertRowsOfOneDeviceStatement.java | 42 ++++ .../plan/statement/crud/InsertRowsStatement.java | 41 ++++ .../plan/statement/crud/InsertTabletStatement.java | 155 +++++++++++++- .../db/wal/recover/file/TsFilePlanRedoer.java | 4 - .../db/engine/storagegroup/DataRegionTest.java | 16 +- .../org/apache/iotdb/db/wal/io/WALFileTest.java | 25 +-- .../iotdb/db/wal/node/ConsensusReqReaderTest.java | 27 +-- .../org/apache/iotdb/db/wal/node/WALNodeTest.java | 26 ++- .../db/wal/recover/file/TsFilePlanRedoerTest.java | 32 +-- .../file/UnsealedTsFileRecoverPerformerTest.java | 1 + 27 files changed, 853 insertions(+), 615 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/RegionWriteExecutor.java b/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/RegionWriteExecutor.java index c035edabe5f..6c77e2be185 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/RegionWriteExecutor.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/execution/executor/RegionWriteExecutor.java @@ -24,7 +24,6 @@ import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.consensus.ConsensusGroupId; import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.commons.consensus.SchemaRegionId; -import org.apache.iotdb.commons.exception.IoTDBException; import org.apache.iotdb.commons.exception.MetadataException; import org.apache.iotdb.commons.path.MeasurementPath; import org.apache.iotdb.commons.path.PartialPath; @@ -36,12 +35,10 @@ import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.consensus.DataRegionConsensusImpl; import org.apache.iotdb.db.consensus.SchemaRegionConsensusImpl; import org.apache.iotdb.db.exception.metadata.MeasurementAlreadyExistException; -import org.apache.iotdb.db.exception.sql.SemanticException; import org.apache.iotdb.db.metadata.schemaregion.ISchemaRegion; import org.apache.iotdb.db.metadata.schemaregion.SchemaEngine; import org.apache.iotdb.db.metadata.template.ClusterTemplateManager; import org.apache.iotdb.db.metadata.template.Template; -import org.apache.iotdb.db.mpp.plan.analyze.schema.SchemaValidator; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor; import org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.ActivateTemplateNode; @@ -224,36 +221,8 @@ public class RegionWriteExecutor { private RegionExecutionResult executeDataInsert( InsertNode insertNode, WritePlanNodeExecutionContext context) { RegionExecutionResult response = new RegionExecutionResult(); - // data insertion should be blocked by data deletion, especially when deleting timeseries - final long startTime = System.nanoTime(); context.getRegionWriteValidationRWLock().readLock().lock(); try { - try { - SchemaValidator.validate(insertNode); - } catch (SemanticException e) { - response.setAccepted(false); - response.setMessage(e.getMessage()); - if (e.getCause() instanceof IoTDBException) { - IoTDBException ioTDBException = (IoTDBException) e.getCause(); - response.setStatus( - RpcUtils.getStatus(ioTDBException.getErrorCode(), ioTDBException.getMessage())); - } else { - response.setStatus(RpcUtils.getStatus(TSStatusCode.METADATA_ERROR, e.getMessage())); - } - return response; - } finally { - PERFORMANCE_OVERVIEW_METRICS.recordScheduleSchemaValidateCost( - System.nanoTime() - startTime); - } - boolean hasFailedMeasurement = insertNode.hasFailedMeasurements(); - String partialInsertMessage = null; - if (hasFailedMeasurement) { - partialInsertMessage = - String.format( - "Fail to insert measurements %s caused by %s", - insertNode.getFailedMeasurements(), insertNode.getFailedMessages()); - LOGGER.warn(partialInsertMessage); - } ConsensusWriteResponse writeResponse = fireTriggerAndInsert(context.getRegionId(), insertNode); @@ -261,17 +230,10 @@ public class RegionWriteExecutor { // TODO need consider more status if (writeResponse.getStatus() != null) { response.setAccepted( - !hasFailedMeasurement - && TSStatusCode.SUCCESS_STATUS.getStatusCode() - == writeResponse.getStatus().getCode()); + TSStatusCode.SUCCESS_STATUS.getStatusCode() == writeResponse.getStatus().getCode()); if (TSStatusCode.SUCCESS_STATUS.getStatusCode() != writeResponse.getStatus().getCode()) { response.setMessage(writeResponse.getStatus().message); response.setStatus(writeResponse.getStatus()); - } else if (hasFailedMeasurement) { - response.setMessage(partialInsertMessage); - response.setStatus( - RpcUtils.getStatus( - TSStatusCode.METADATA_ERROR.getStatusCode(), partialInsertMessage)); } else { response.setMessage(writeResponse.getStatus().message); } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analysis.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analysis.java index 6b3f4780a67..3ff74d63987 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analysis.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/Analysis.java @@ -22,6 +22,7 @@ package org.apache.iotdb.db.mpp.plan.analyze; import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation; import org.apache.iotdb.common.rpc.thrift.TEndPoint; import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; +import org.apache.iotdb.common.rpc.thrift.TSStatus; import org.apache.iotdb.common.rpc.thrift.TSchemaNode; import org.apache.iotdb.commons.partition.DataPartition; import org.apache.iotdb.commons.partition.SchemaPartition; @@ -80,9 +81,10 @@ public class Analysis { private boolean finishQueryAfterAnalyze; - // potential fail message when finishQueryAfterAnalyze is true. If failMessage is NULL, means no + // potential fail status when finishQueryAfterAnalyze is true. If failStatus is NULL, means no // fail. - private String failMessage; + + private TSStatus failStatus; ///////////////////////////////////////////////////////////////////////////////////////////////// // Query Analysis (used in ALIGN BY TIME) @@ -407,15 +409,15 @@ public class Analysis { } public boolean isFailed() { - return failMessage != null; + return failStatus != null; } - public String getFailMessage() { - return failMessage; + public TSStatus getFailStatus() { + return this.failStatus; } - public void setFailMessage(String failMessage) { - this.failMessage = failMessage; + public void setFailStatus(TSStatus status) { + this.failStatus = status; } public void setDeviceViewInputIndexesMap(Map<String, List<Integer>> deviceViewInputIndexesMap) { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java index 30162d3ad06..f30385724bf 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/AnalyzeVisitor.java @@ -136,6 +136,7 @@ import org.apache.iotdb.db.mpp.plan.statement.sys.sync.ShowPipeSinkTypeStatement import org.apache.iotdb.db.query.control.SessionManager; import org.apache.iotdb.db.utils.FileLoaderUtils; import org.apache.iotdb.db.utils.TimePartitionUtils; +import org.apache.iotdb.rpc.RpcUtils; import org.apache.iotdb.rpc.TSStatusCode; import org.apache.iotdb.tsfile.common.constant.TsFileConstant; import org.apache.iotdb.tsfile.file.metadata.TimeseriesMetadata; @@ -2025,32 +2026,48 @@ public class AnalyzeVisitor extends StatementVisitor<Analysis, MPPQueryContext> public Analysis visitInsertTablet( InsertTabletStatement insertTabletStatement, MPPQueryContext context) { context.setQueryType(QueryType.WRITE); + Analysis analysis = new Analysis(); + analysis.setStatement(insertTabletStatement); + insertTabletStatement.validateSchema(analysis); + if (analysis.isFinishQueryAfterAnalyze()) { + return analysis; + } DataPartitionQueryParam dataPartitionQueryParam = new DataPartitionQueryParam(); dataPartitionQueryParam.setDevicePath(insertTabletStatement.getDevicePath().getFullPath()); dataPartitionQueryParam.setTimePartitionSlotList(insertTabletStatement.getTimePartitionSlots()); - return getAnalysisForWriting( - insertTabletStatement, Collections.singletonList(dataPartitionQueryParam)); + return getAnalysisForWriting(analysis, Collections.singletonList(dataPartitionQueryParam)); } @Override public Analysis visitInsertRow(InsertRowStatement insertRowStatement, MPPQueryContext context) { context.setQueryType(QueryType.WRITE); + Analysis analysis = new Analysis(); + analysis.setStatement(insertRowStatement); + insertRowStatement.validateSchema(analysis); + if (analysis.isFinishQueryAfterAnalyze()) { + return analysis; + } DataPartitionQueryParam dataPartitionQueryParam = new DataPartitionQueryParam(); dataPartitionQueryParam.setDevicePath(insertRowStatement.getDevicePath().getFullPath()); dataPartitionQueryParam.setTimePartitionSlotList( Collections.singletonList(insertRowStatement.getTimePartitionSlot())); - return getAnalysisForWriting( - insertRowStatement, Collections.singletonList(dataPartitionQueryParam)); + return getAnalysisForWriting(analysis, Collections.singletonList(dataPartitionQueryParam)); } @Override public Analysis visitInsertRows( InsertRowsStatement insertRowsStatement, MPPQueryContext context) { context.setQueryType(QueryType.WRITE); + Analysis analysis = new Analysis(); + analysis.setStatement(insertRowsStatement); + insertRowsStatement.validateSchema(analysis); + if (analysis.isFinishQueryAfterAnalyze()) { + return analysis; + } Map<String, Set<TTimePartitionSlot>> dataPartitionQueryParamMap = new HashMap<>(); for (InsertRowStatement insertRowStatement : insertRowsStatement.getInsertRowStatementList()) { @@ -2068,13 +2085,19 @@ public class AnalyzeVisitor extends StatementVisitor<Analysis, MPPQueryContext> dataPartitionQueryParams.add(dataPartitionQueryParam); } - return getAnalysisForWriting(insertRowsStatement, dataPartitionQueryParams); + return getAnalysisForWriting(analysis, dataPartitionQueryParams); } @Override public Analysis visitInsertMultiTablets( InsertMultiTabletsStatement insertMultiTabletsStatement, MPPQueryContext context) { context.setQueryType(QueryType.WRITE); + Analysis analysis = new Analysis(); + analysis.setStatement(insertMultiTabletsStatement); + insertMultiTabletsStatement.validateSchema(analysis); + if (analysis.isFinishQueryAfterAnalyze()) { + return analysis; + } Map<String, Set<TTimePartitionSlot>> dataPartitionQueryParamMap = new HashMap<>(); for (InsertTabletStatement insertTabletStatement : @@ -2093,13 +2116,19 @@ public class AnalyzeVisitor extends StatementVisitor<Analysis, MPPQueryContext> dataPartitionQueryParams.add(dataPartitionQueryParam); } - return getAnalysisForWriting(insertMultiTabletsStatement, dataPartitionQueryParams); + return getAnalysisForWriting(analysis, dataPartitionQueryParams); } @Override public Analysis visitInsertRowsOfOneDevice( InsertRowsOfOneDeviceStatement insertRowsOfOneDeviceStatement, MPPQueryContext context) { context.setQueryType(QueryType.WRITE); + Analysis analysis = new Analysis(); + analysis.setStatement(insertRowsOfOneDeviceStatement); + insertRowsOfOneDeviceStatement.validateSchema(analysis); + if (analysis.isFinishQueryAfterAnalyze()) { + return analysis; + } DataPartitionQueryParam dataPartitionQueryParam = new DataPartitionQueryParam(); dataPartitionQueryParam.setDevicePath( @@ -2107,8 +2136,7 @@ public class AnalyzeVisitor extends StatementVisitor<Analysis, MPPQueryContext> dataPartitionQueryParam.setTimePartitionSlotList( insertRowsOfOneDeviceStatement.getTimePartitionSlots()); - return getAnalysisForWriting( - insertRowsOfOneDeviceStatement, Collections.singletonList(dataPartitionQueryParam)); + return getAnalysisForWriting(analysis, Collections.singletonList(dataPartitionQueryParam)); } @Override @@ -2193,16 +2221,16 @@ public class AnalyzeVisitor extends StatementVisitor<Analysis, MPPQueryContext> /** get analysis according to statement and params */ private Analysis getAnalysisForWriting( - Statement statement, List<DataPartitionQueryParam> dataPartitionQueryParams) { - Analysis analysis = new Analysis(); - analysis.setStatement(statement); + Analysis analysis, List<DataPartitionQueryParam> dataPartitionQueryParams) { DataPartition dataPartition = partitionFetcher.getOrCreateDataPartition(dataPartitionQueryParams); if (dataPartition.isEmpty()) { analysis.setFinishQueryAfterAnalyze(true); - analysis.setFailMessage( - "Database not exists and failed to create automatically because enable_auto_create_schema is FALSE."); + analysis.setFailStatus( + RpcUtils.getStatus( + TSStatusCode.DATABASE_NOT_EXIST.getStatusCode(), + "Database not exists and failed to create automatically because enable_auto_create_schema is FALSE.")); } analysis.setDataPartitionInfo(dataPartition); return analysis; diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/schema/SchemaValidator.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/schema/SchemaValidator.java index 6cbd587a00b..1c8655843c9 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/schema/SchemaValidator.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/analyze/schema/SchemaValidator.java @@ -23,8 +23,10 @@ import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.exception.query.QueryProcessException; import org.apache.iotdb.db.exception.sql.SemanticException; import org.apache.iotdb.db.mpp.common.schematree.ISchemaTree; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.BatchInsertNode; -import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertNode; +import org.apache.iotdb.db.mpp.plan.statement.crud.InsertBaseStatement; +import org.apache.iotdb.db.mpp.plan.statement.crud.InsertMultiTabletsStatement; +import org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowsOfOneDeviceStatement; +import org.apache.iotdb.db.mpp.plan.statement.crud.InsertRowsStatement; 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; @@ -35,17 +37,19 @@ public class SchemaValidator { private static final ISchemaFetcher SCHEMA_FETCHER = ClusterSchemaFetcher.getInstance(); - public static void validate(InsertNode insertNode) { + public static void validate(InsertBaseStatement insertStatement) { try { - if (insertNode instanceof BatchInsertNode) { + if (insertStatement instanceof InsertRowsStatement + || insertStatement instanceof InsertMultiTabletsStatement + || insertStatement instanceof InsertRowsOfOneDeviceStatement) { SCHEMA_FETCHER.fetchAndComputeSchemaWithAutoCreate( - ((BatchInsertNode) insertNode).getSchemaValidationList()); + insertStatement.getSchemaValidationList()); } else { - SCHEMA_FETCHER.fetchAndComputeSchemaWithAutoCreate(insertNode.getSchemaValidation()); + SCHEMA_FETCHER.fetchAndComputeSchemaWithAutoCreate(insertStatement.getSchemaValidation()); } - insertNode.updateAfterSchemaValidation(); + insertStatement.updateAfterSchemaValidation(); } catch (QueryProcessException e) { - throw new SemanticException(e); + throw new SemanticException(e.getMessage()); } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java index 715cf29c716..d34c31bea8f 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/QueryExecution.java @@ -202,7 +202,7 @@ public class QueryExecution implements IQueryExecution { if (skipExecute()) { logger.debug("[SkipExecute]"); if (context.getQueryType() == QueryType.WRITE && analysis.isFailed()) { - stateMachine.transitionToFailed(new RuntimeException(analysis.getFailMessage())); + stateMachine.transitionToFailed(analysis.getFailStatus()); } else { constructResultForMemorySource(); stateMachine.transitionToRunning(); @@ -224,6 +224,11 @@ public class QueryExecution implements IQueryExecution { } PERFORMANCE_OVERVIEW_METRICS.recordPlanCost(System.nanoTime() - startTime); schedule(); + + // set partial insert error message + if (context.getQueryType() == QueryType.WRITE && analysis.isFailed()) { + stateMachine.transitionToFailed(analysis.getFailStatus()); + } } private void checkTimeOutForQuery() { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java index a8cb30fac97..bc6065d039c 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java @@ -447,6 +447,7 @@ public class LogicalPlanVisitor extends StatementVisitor<PlanNode, MPPQueryConte insertTabletStatement.isAligned(), insertTabletStatement.getMeasurements(), insertTabletStatement.getDataTypes(), + insertTabletStatement.getMeasurementSchemas(), insertTabletStatement.getTimes(), insertTabletStatement.getBitMaps(), insertTabletStatement.getColumns(), @@ -462,6 +463,7 @@ public class LogicalPlanVisitor extends StatementVisitor<PlanNode, MPPQueryConte insertRowStatement.isAligned(), insertRowStatement.getMeasurements(), insertRowStatement.getDataTypes(), + insertRowStatement.getMeasurementSchemas(), insertRowStatement.getTime(), insertRowStatement.getValues(), insertRowStatement.isNeedInferType()); @@ -637,6 +639,7 @@ public class LogicalPlanVisitor extends StatementVisitor<PlanNode, MPPQueryConte insertRowStatement.isAligned(), insertRowStatement.getMeasurements(), insertRowStatement.getDataTypes(), + insertRowStatement.getMeasurementSchemas(), insertRowStatement.getTime(), insertRowStatement.getValues(), insertRowStatement.isNeedInferType()), @@ -661,6 +664,7 @@ public class LogicalPlanVisitor extends StatementVisitor<PlanNode, MPPQueryConte insertTabletStatement.isAligned(), insertTabletStatement.getMeasurements(), insertTabletStatement.getDataTypes(), + insertTabletStatement.getMeasurementSchemas(), insertTabletStatement.getTimes(), insertTabletStatement.getBitMaps(), insertTabletStatement.getColumns(), @@ -689,6 +693,7 @@ public class LogicalPlanVisitor extends StatementVisitor<PlanNode, MPPQueryConte insertRowStatement.isAligned(), insertRowStatement.getMeasurements(), insertRowStatement.getDataTypes(), + insertRowStatement.getMeasurementSchemas(), insertRowStatement.getTime(), insertRowStatement.getValues(), insertRowStatement.isNeedInferType())); diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/BatchInsertNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/BatchInsertNode.java deleted file mode 100644 index a0c774d5c49..00000000000 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/BatchInsertNode.java +++ /dev/null @@ -1,33 +0,0 @@ -/* - * 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.write; - -import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; - -import java.util.List; - -/** - * BatchInsertNode contains multiple sub insert. Insert node which contains multiple sub insert - * nodes needs to implement it. - */ -public interface BatchInsertNode { - - List<ISchemaValidation> getSchemaValidationList(); -} diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertMultiTabletsNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertMultiTabletsNode.java index d0b46b103ae..a60e0c4cd19 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertMultiTabletsNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertMultiTabletsNode.java @@ -42,7 +42,7 @@ import java.util.Map; import java.util.Objects; import java.util.stream.Collectors; -public class InsertMultiTabletsNode extends InsertNode implements BatchInsertNode { +public class InsertMultiTabletsNode extends InsertNode { /** * the value is used to indict the parent InsertTabletNode's index when the parent @@ -135,11 +135,6 @@ public class InsertMultiTabletsNode extends InsertNode implements BatchInsertNod insertTabletNodeList.forEach(plan -> plan.setSearchIndex(index)); } - @Override - protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { - return false; - } - @Override public List<WritePlanNode> splitByPartition(Analysis analysis) { Map<TRegionReplicaSet, InsertMultiTabletsNode> splitMap = new HashMap<>(); @@ -193,13 +188,6 @@ public class InsertMultiTabletsNode extends InsertNode implements BatchInsertNod return null; } - @Override - public List<ISchemaValidation> getSchemaValidationList() { - return insertTabletNodeList.stream() - .map(InsertTabletNode::getSchemaValidation) - .collect(Collectors.toList()); - } - public static InsertMultiTabletsNode deserialize(ByteBuffer byteBuffer) { PlanNodeId planNodeId; List<InsertTabletNode> insertTabletNodeList = new ArrayList<>(); @@ -279,8 +267,4 @@ public class InsertMultiTabletsNode extends InsertNode implements BatchInsertNod throw new NotImplementedException(); } - @Override - public Object getFirstValueOfIndex(int index) { - throw new NotImplementedException(); - } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertNode.java index 0994935c958..8eaf854cd72 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertNode.java @@ -20,13 +20,10 @@ package org.apache.iotdb.db.mpp.plan.planner.plan.node.write; import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.commons.utils.TestOnly; import org.apache.iotdb.consensus.iot.wal.ConsensusReqReader; import org.apache.iotdb.db.conf.IoTDBDescriptor; -import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException; -import org.apache.iotdb.db.exception.metadata.PathNotExistException; -import org.apache.iotdb.db.exception.query.QueryProcessException; import org.apache.iotdb.db.metadata.idtable.entry.IDeviceID; -import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.plan.planner.plan.node.WritePlanNode; import org.apache.iotdb.db.wal.buffer.IWALByteBufferView; @@ -41,11 +38,7 @@ import java.io.DataOutputStream; import java.io.IOException; import java.nio.ByteBuffer; import java.util.Arrays; -import java.util.Collections; -import java.util.List; -import java.util.Map; import java.util.Objects; -import java.util.stream.Collectors; public abstract class InsertNode extends WritePlanNode { @@ -62,11 +55,8 @@ public abstract class InsertNode extends WritePlanNode { protected MeasurementSchema[] measurementSchemas; protected String[] measurements; protected TSDataType[] dataTypes; - // TODO(INSERT) need to change it to a function handle to update last time value - // protected IMeasurementMNode[] measurementMNodes; - /** index of failed measurements -> info including measurement, data type and value */ - protected Map<Integer, FailedMeasurementInfo> failedMeasurementIndex2Info; + protected int failedMeasurementNumber = 0; /** * device id reference, for reuse device id in both id table and memtable <br> @@ -242,66 +232,10 @@ public abstract class InsertNode extends WritePlanNode { return dataRegionReplicaSet; } - public ISchemaValidation getSchemaValidation() { - throw new UnsupportedOperationException(); - } - - public void updateAfterSchemaValidation() throws QueryProcessException {} - - /** Check whether data types are matched with measurement schemas */ - protected void selfCheckDataTypes(int index) - throws DataTypeMismatchException, PathNotExistException { - if (IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert()) { - // if enable partial insert, mark failed measurements with exception - if (measurementSchemas[index] == null) { - markFailedMeasurement( - index, - new PathNotExistException(devicePath.concatNode(measurements[index]).getFullPath())); - } else if ((dataTypes[index] != measurementSchemas[index].getType() - && !checkAndCastDataType(index, measurementSchemas[index].getType()))) { - markFailedMeasurement( - index, - new DataTypeMismatchException( - devicePath.getFullPath(), - measurements[index], - dataTypes[index], - measurementSchemas[index].getType(), - getMinTime(), - getFirstValueOfIndex(index))); - } - } else { - // if not enable partial insert, throw the exception directly - if (measurementSchemas[index] == null) { - throw new PathNotExistException(devicePath.concatNode(measurements[index]).getFullPath()); - } else if ((dataTypes[index] != measurementSchemas[index].getType() - && !checkAndCastDataType(index, measurementSchemas[index].getType()))) { - throw new DataTypeMismatchException( - devicePath.getFullPath(), - measurements[index], - dataTypes[index], - measurementSchemas[index].getType(), - getMinTime(), - getFirstValueOfIndex(index)); - } - } - } - - protected abstract boolean checkAndCastDataType(int columnIndex, TSDataType dataType); - public abstract long getMinTime(); - public abstract Object getFirstValueOfIndex(int index); - // region partial insert - /** - * Mark failed measurement, measurements[index], dataTypes[index] and values/columns[index] would - * be null. We'd better use "measurements[index] == null" to determine if the measurement failed. - * <br> - * This method is not concurrency-safe. - * - * @param index failed measurement index - * @param cause cause Exception of failure - */ + @TestOnly public void markFailedMeasurement(int index, Exception cause) { throw new UnsupportedOperationException(); } @@ -316,57 +250,20 @@ public abstract class InsertNode extends WritePlanNode { } public boolean hasFailedMeasurements() { - return failedMeasurementIndex2Info != null && !failedMeasurementIndex2Info.isEmpty(); + for (Object o : measurements) { + if (o == null) { + failedMeasurementNumber++; + } + } + return failedMeasurementNumber != 0; + } + + public void setFailedMeasurementNumber(int failedMeasurementNumber) { + this.failedMeasurementNumber = failedMeasurementNumber; } public int getFailedMeasurementNumber() { - return failedMeasurementIndex2Info == null ? 0 : failedMeasurementIndex2Info.size(); - } - - public List<String> getFailedMeasurements() { - return failedMeasurementIndex2Info == null - ? Collections.emptyList() - : failedMeasurementIndex2Info.values().stream() - .map(info -> info.measurement) - .collect(Collectors.toList()); - } - - public List<Exception> getFailedExceptions() { - return failedMeasurementIndex2Info == null - ? Collections.emptyList() - : failedMeasurementIndex2Info.values().stream() - .map(info -> info.cause) - .collect(Collectors.toList()); - } - - public List<String> getFailedMessages() { - return failedMeasurementIndex2Info == null - ? Collections.emptyList() - : failedMeasurementIndex2Info.values().stream() - .map( - info -> { - Throwable cause = info.cause; - while (cause.getCause() != null) { - cause = cause.getCause(); - } - return cause.getMessage(); - }) - .collect(Collectors.toList()); - } - - protected static class FailedMeasurementInfo { - protected String measurement; - protected TSDataType dataType; - protected Object value; - protected Exception cause; - - public FailedMeasurementInfo( - String measurement, TSDataType dataType, Object value, Exception cause) { - this.measurement = measurement; - this.dataType = dataType; - this.value = value; - this.cause = cause; - } + return failedMeasurementNumber; } // endregion diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowNode.java index 311dfd28463..5b0472d6504 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowNode.java @@ -66,7 +66,7 @@ import java.util.HashMap; import java.util.List; import java.util.Objects; -public class InsertRowNode extends InsertNode implements WALEntryValue, ISchemaValidation { +public class InsertRowNode extends InsertNode implements WALEntryValue { private static final Logger logger = LoggerFactory.getLogger(InsertRowNode.class); @@ -98,6 +98,23 @@ public class InsertRowNode extends InsertNode implements WALEntryValue, ISchemaV this.isNeedInferType = isNeedInferType; } + public InsertRowNode( + PlanNodeId id, + PartialPath devicePath, + boolean isAligned, + String[] measurements, + TSDataType[] dataTypes, + MeasurementSchema[] measurementSchemas, + long time, + Object[] values, + boolean isNeedInferType) { + super(id, devicePath, isAligned, measurements, dataTypes); + this.measurementSchemas = measurementSchemas; + this.time = time; + this.values = values; + this.isNeedInferType = isNeedInferType; + } + @Override public List<WritePlanNode> splitByPartition(Analysis analysis) { TTimePartitionSlot timePartitionSlot = TimePartitionUtils.getTimePartition(time); @@ -157,16 +174,6 @@ public class InsertRowNode extends InsertNode implements WALEntryValue, ISchemaV } } - @Override - public TSEncoding getEncoding(int index) { - return null; - } - - @Override - public CompressionType getCompressionType(int index) { - return null; - } - public Object[] getValues() { return values; } @@ -196,82 +203,11 @@ public class InsertRowNode extends InsertNode implements WALEntryValue, ISchemaV return Collections.singletonList(TimePartitionUtils.getTimePartition(time)); } - @Override - protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { - if (CommonUtils.checkCanCastType(dataTypes[columnIndex], dataType)) { - logger.warn( - "Inserting to {}.{} : Cast from {} to {}", - devicePath, - measurements[columnIndex], - dataTypes[columnIndex], - dataType); - values[columnIndex] = - CommonUtils.castValue(dataTypes[columnIndex], dataType, values[columnIndex]); - dataTypes[columnIndex] = dataType; - return true; - } - return false; - } - - /** - * transfer String[] values to specific data types when isNeedInferType is true. <br> - * Notice: measurementSchemas must be initialized before calling this method - */ - @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity warning - public void transferType() throws QueryProcessException { - - for (int i = 0; i < measurementSchemas.length; i++) { - // null when time series doesn't exist - if (measurementSchemas[i] == null) { - if (!IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert()) { - throw new QueryProcessException( - new PathNotExistException( - devicePath.getFullPath() + IoTDBConstant.PATH_SEPARATOR + measurements[i])); - } else { - markFailedMeasurement( - i, - new QueryProcessException( - new PathNotExistException( - devicePath.getFullPath() + IoTDBConstant.PATH_SEPARATOR + measurements[i]))); - } - continue; - } - // parse string value to specific type - dataTypes[i] = measurementSchemas[i].getType(); - try { - values[i] = CommonUtils.parseValue(dataTypes[i], values[i].toString()); - } catch (Exception e) { - logger.warn( - "data type of {}.{} is not consistent, registered type {}, inserting timestamp {}, value {}", - devicePath, - measurements[i], - dataTypes[i], - time, - values[i]); - if (!IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert()) { - throw e; - } else { - markFailedMeasurement(i, e); - } - } - } - isNeedInferType = false; - } - @Override public void markFailedMeasurement(int index, Exception cause) { if (measurements[index] == null) { return; } - - if (failedMeasurementIndex2Info == null) { - failedMeasurementIndex2Info = new HashMap<>(); - } - - FailedMeasurementInfo failedMeasurementInfo = - new FailedMeasurementInfo(measurements[index], dataTypes[index], values[index], cause); - failedMeasurementIndex2Info.putIfAbsent(index, failedMeasurementInfo); - measurements[index] = null; dataTypes[index] = null; values[index] = null; @@ -527,11 +463,6 @@ public class InsertRowNode extends InsertNode implements WALEntryValue, ISchemaV return getTime(); } - @Override - public Object getFirstValueOfIndex(int index) { - return values[index]; - } - // region serialize & deserialize methods for WAL /** Serialized size for wal */ @Override @@ -810,49 +741,4 @@ public class InsertRowNode extends InsertNode implements WALEntryValue, ISchemaV Object value = values[columnIndex]; return new TimeValuePair(time, TsPrimitiveType.getByType(dataTypes[columnIndex], value)); } - - @Override - public void validateDeviceSchema(boolean isAligned) { - if (this.isAligned != isAligned) { - throw new SemanticException( - new AlignedTimeseriesException( - String.format( - "timeseries under this device are%s aligned, " + "please use %s interface", - isAligned ? "" : " not", isAligned ? "aligned" : "non-aligned"), - devicePath.getFullPath())); - } - } - - @Override - public ISchemaValidation getSchemaValidation() { - return this; - } - - @Override - public void updateAfterSchemaValidation() throws QueryProcessException { - if (isNeedInferType) { - transferType(); - } - } - - @Override - public void validateMeasurementSchema(int index, IMeasurementSchemaInfo measurementSchemaInfo) { - if (measurementSchemas == null) { - measurementSchemas = new MeasurementSchema[measurements.length]; - } - if (measurementSchemaInfo == null) { - measurementSchemas[index] = null; - } else { - measurementSchemas[index] = measurementSchemaInfo.getSchema(); - } - if (isNeedInferType) { - return; - } - - try { - selfCheckDataTypes(index); - } catch (DataTypeMismatchException | PathNotExistException e) { - throw new SemanticException(e); - } - } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsNode.java index 5921c033f6b..c043e5991b9 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsNode.java @@ -45,7 +45,7 @@ import java.util.Map; import java.util.Objects; import java.util.stream.Collectors; -public class InsertRowsNode extends InsertNode implements BatchInsertNode { +public class InsertRowsNode extends InsertNode { /** * Suppose there is an InsertRowsNode, which contains 5 InsertRowNodes, @@ -121,11 +121,6 @@ public class InsertRowsNode extends InsertNode implements BatchInsertNode { @Override public void addChild(PlanNode child) {} - @Override - protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { - return false; - } - @Override public boolean equals(Object o) { if (this == o) return true; @@ -156,23 +151,6 @@ public class InsertRowsNode extends InsertNode implements BatchInsertNode { return null; } - @Override - public List<ISchemaValidation> getSchemaValidationList() { - return insertRowNodeList.stream() - .map(InsertRowNode::getSchemaValidation) - .collect(Collectors.toList()); - } - - @Override - public void updateAfterSchemaValidation() throws QueryProcessException { - for (InsertRowNode insertRowNode : insertRowNodeList) { - insertRowNode.updateAfterSchemaValidation(); - if (!this.hasFailedMeasurements() && insertRowNode.hasFailedMeasurements()) { - this.failedMeasurementIndex2Info = insertRowNode.failedMeasurementIndex2Info; - } - } - } - public static InsertRowsNode deserialize(ByteBuffer byteBuffer) { PlanNodeId planNodeId; List<InsertRowNode> insertRowNodeList = new ArrayList<>(); @@ -267,8 +245,4 @@ public class InsertRowsNode extends InsertNode implements BatchInsertNode { throw new NotImplementedException(); } - @Override - public Object getFirstValueOfIndex(int index) { - throw new NotImplementedException(); - } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java index 32b3e1c08d2..4f101942688 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertRowsOfOneDeviceNode.java @@ -49,7 +49,7 @@ import java.util.Objects; import java.util.Set; import java.util.stream.Collectors; -public class InsertRowsOfOneDeviceNode extends InsertNode implements BatchInsertNode { +public class InsertRowsOfOneDeviceNode extends InsertNode { /** * Suppose there is an InsertRowsOfOneDeviceNode, which contains 5 InsertRowNodes, @@ -144,11 +144,6 @@ public class InsertRowsOfOneDeviceNode extends InsertNode implements BatchInsert return null; } - @Override - protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { - return false; - } - @Override public List<WritePlanNode> splitByPartition(Analysis analysis) { List<WritePlanNode> result = new ArrayList<>(); @@ -291,23 +286,6 @@ public class InsertRowsOfOneDeviceNode extends InsertNode implements BatchInsert return Objects.hash(super.hashCode(), insertRowNodeIndexList, insertRowNodeList); } - @Override - public List<ISchemaValidation> getSchemaValidationList() { - return insertRowNodeList.stream() - .map(InsertRowNode::getSchemaValidation) - .collect(Collectors.toList()); - } - - @Override - public void updateAfterSchemaValidation() throws QueryProcessException { - for (InsertRowNode insertRowNode : insertRowNodeList) { - insertRowNode.updateAfterSchemaValidation(); - if (!this.hasFailedMeasurements() && insertRowNode.hasFailedMeasurements()) { - this.failedMeasurementIndex2Info = insertRowNode.failedMeasurementIndex2Info; - } - } - } - @Override public <R, C> R accept(PlanVisitor<R, C> visitor, C context) { return visitor.visitInsertRowsOfOneDevice(this, context); @@ -317,9 +295,4 @@ public class InsertRowsOfOneDeviceNode extends InsertNode implements BatchInsert public long getMinTime() { throw new NotImplementedException(); } - - @Override - public Object getFirstValueOfIndex(int index) { - throw new NotImplementedException(); - } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertTabletNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertTabletNode.java index b292b4f49c6..b60f4ffb6a3 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertTabletNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/InsertTabletNode.java @@ -23,13 +23,7 @@ import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.commons.utils.TestOnly; -import org.apache.iotdb.db.exception.metadata.AlignedTimeseriesException; -import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException; -import org.apache.iotdb.db.exception.metadata.PathNotExistException; -import org.apache.iotdb.db.exception.sql.SemanticException; -import org.apache.iotdb.db.mpp.common.schematree.IMeasurementSchemaInfo; import org.apache.iotdb.db.mpp.plan.analyze.Analysis; -import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; 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; @@ -43,9 +37,7 @@ import org.apache.iotdb.db.wal.buffer.WALEntryValue; import org.apache.iotdb.db.wal.utils.WALWriteUtils; import org.apache.iotdb.tsfile.exception.NotImplementedException; import org.apache.iotdb.tsfile.exception.write.UnSupportedDataTypeException; -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.apache.iotdb.tsfile.read.TimeValuePair; import org.apache.iotdb.tsfile.utils.Binary; import org.apache.iotdb.tsfile.utils.BitMap; @@ -69,7 +61,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; -public class InsertTabletNode extends InsertNode implements WALEntryValue, ISchemaValidation { +public class InsertTabletNode extends InsertNode implements WALEntryValue { private static final Logger logger = LoggerFactory.getLogger(InsertTabletNode.class); @@ -103,11 +95,13 @@ public class InsertTabletNode extends InsertNode implements WALEntryValue, ISche boolean isAligned, String[] measurements, TSDataType[] dataTypes, + MeasurementSchema[] measurementSchemas, long[] times, BitMap[] bitMaps, Object[] columns, int rowCount) { super(id, devicePath, isAligned, measurements, dataTypes); + this.measurementSchemas = measurementSchemas; this.times = times; this.bitMaps = bitMaps; this.columns = columns; @@ -177,23 +171,6 @@ public class InsertTabletNode extends InsertNode implements WALEntryValue, ISche return null; } - @Override - protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { - if (CommonUtils.checkCanCastType(dataTypes[columnIndex], dataType)) { - logger.warn( - "Inserting to {}.{} : Cast from {} to {}", - devicePath, - measurements[columnIndex], - dataTypes[columnIndex], - dataType); - columns[columnIndex] = - CommonUtils.castArray(dataTypes[columnIndex], dataType, columns[columnIndex]); - dataTypes[columnIndex] = dataType; - return true; - } - return false; - } - @Override public List<WritePlanNode> splitByPartition(Analysis analysis) { // only single device in single database @@ -287,6 +264,7 @@ public class InsertTabletNode extends InsertNode implements WALEntryValue, ISche isAligned, measurements, dataTypes, + measurementSchemas, subTimes, bitMaps, values, @@ -361,15 +339,6 @@ public class InsertTabletNode extends InsertNode implements WALEntryValue, ISche if (measurements[index] == null) { return; } - - if (failedMeasurementIndex2Info == null) { - failedMeasurementIndex2Info = new HashMap<>(); - } - - FailedMeasurementInfo failedMeasurementInfo = - new FailedMeasurementInfo(measurements[index], dataTypes[index], columns[index], cause); - failedMeasurementIndex2Info.putIfAbsent(index, failedMeasurementInfo); - measurements[index] = null; dataTypes[index] = null; columns[index] = null; @@ -380,41 +349,6 @@ public class InsertTabletNode extends InsertNode implements WALEntryValue, ISche return times[0]; } - @Override - public Object getFirstValueOfIndex(int index) { - Object value; - switch (dataTypes[index]) { - case INT32: - int[] intValues = (int[]) columns[index]; - value = intValues[0]; - break; - case INT64: - long[] longValues = (long[]) columns[index]; - value = longValues[0]; - break; - case FLOAT: - float[] floatValues = (float[]) columns[index]; - value = floatValues[0]; - break; - case DOUBLE: - double[] doubleValues = (double[]) columns[index]; - value = doubleValues[0]; - break; - case BOOLEAN: - boolean[] boolValues = (boolean[]) columns[index]; - value = boolValues[0]; - break; - case TEXT: - Binary[] binaryValues = (Binary[]) columns[index]; - value = binaryValues[0]; - break; - default: - throw new UnSupportedDataTypeException( - String.format(DATATYPE_UNSUPPORTED, dataTypes[index])); - } - return value; - } - @Override protected void serializeAttributes(ByteBuffer byteBuffer) { PlanNodeType.INSERT_TABLET.serialize(byteBuffer); @@ -1122,49 +1056,4 @@ public class InsertTabletNode extends InsertNode implements WALEntryValue, ISche } return new TimeValuePair(times[lastIdx], value); } - - @Override - public TSEncoding getEncoding(int index) { - return null; - } - - @Override - public CompressionType getCompressionType(int index) { - return null; - } - - @Override - public void validateDeviceSchema(boolean isAligned) { - if (this.isAligned != isAligned) { - throw new SemanticException( - new AlignedTimeseriesException( - String.format( - "timeseries under this device are%s aligned, " + "please use %s interface", - isAligned ? "" : " not", isAligned ? "aligned" : "non-aligned"), - devicePath.getFullPath())); - } - } - - @Override - public ISchemaValidation getSchemaValidation() { - return this; - } - - @Override - public void validateMeasurementSchema(int index, IMeasurementSchemaInfo measurementSchemaInfo) { - if (measurementSchemas == null) { - measurementSchemas = new MeasurementSchema[measurements.length]; - } - if (measurementSchemaInfo == null) { - measurementSchemas[index] = null; - } else { - measurementSchemas[index] = measurementSchemaInfo.getSchema(); - } - - try { - selfCheckDataTypes(index); - } catch (DataTypeMismatchException | PathNotExistException e) { - throw new SemanticException(e); - } - } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java index a471e9256aa..c6db146a2c0 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/FragmentInstanceDispatcherImpl.java @@ -37,6 +37,7 @@ import org.apache.iotdb.db.mpp.metric.QueryMetricsManager; import org.apache.iotdb.db.mpp.plan.analyze.QueryType; import org.apache.iotdb.db.mpp.plan.planner.plan.FragmentInstance; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode; +import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertNode; import org.apache.iotdb.db.utils.SetThreadName; import org.apache.iotdb.mpp.rpc.thrift.TFragmentInstance; import org.apache.iotdb.mpp.rpc.thrift.TPlanNode; @@ -131,6 +132,7 @@ public class FragmentInstanceDispatcherImpl implements IFragInstanceDispatcher { private Future<FragInstanceDispatchResult> dispatchWriteSync(List<FragmentInstance> instances) { List<TSStatus> failureStatusList = new ArrayList<>(); for (FragmentInstance instance : instances) { + // TODO:(hhn) try (SetThreadName threadName = new SetThreadName(instance.getId().getFullId())) { dispatchOneInstance(instance); } catch (FragmentInstanceDispatchException e) { @@ -164,10 +166,15 @@ public class FragmentInstanceDispatcherImpl implements IFragInstanceDispatcher { } private Future<FragInstanceDispatchResult> dispatchWriteAsync(List<FragmentInstance> instances) { + List<TSStatus> dataNodeFailureList = new ArrayList<>(); // split local and remote instances List<FragmentInstance> localInstances = new ArrayList<>(); List<FragmentInstance> remoteInstances = new ArrayList<>(); for (FragmentInstance instance : instances) { + PlanNode planNode = instance.getFragment().getPlanNodeTree(); + // TODO:(hhn) + if (planNode instanceof InsertNode) {} + TEndPoint endPoint = instance.getHostDataNode().getInternalEndPoint(); if (isDispatchedToLocal(endPoint)) { localInstances.add(instance); @@ -180,14 +187,12 @@ public class FragmentInstanceDispatcherImpl implements IFragInstanceDispatcher { new AsyncPlanNodeSender(asyncInternalServiceClientManager, remoteInstances); asyncPlanNodeSender.sendAll(); - List<TSStatus> dataNodeFailureList = new ArrayList<>(); - if (!localInstances.isEmpty()) { // sync dispatch to local long localScheduleStartTime = System.nanoTime(); for (FragmentInstance localInstance : localInstances) { try (SetThreadName threadName = new SetThreadName(localInstance.getId().getFullId())) { - dispatchOneInstance(localInstance); + dispatchLocally(localInstance); } catch (FragmentInstanceDispatchException e) { dataNodeFailureList.add(e.getFailureStatus()); } catch (Throwable t) { diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertBaseStatement.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertBaseStatement.java index d2d6d55ea12..27e790ea74a 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertBaseStatement.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertBaseStatement.java @@ -18,15 +18,40 @@ */ package org.apache.iotdb.db.mpp.plan.statement.crud; +import org.apache.iotdb.commons.exception.IoTDBException; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.commons.service.metric.enums.PerformanceOverviewMetrics; +import org.apache.iotdb.db.conf.IoTDBDescriptor; +import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException; +import org.apache.iotdb.db.exception.metadata.PathNotExistException; +import org.apache.iotdb.db.exception.query.QueryProcessException; +import org.apache.iotdb.db.exception.sql.SemanticException; +import org.apache.iotdb.db.mpp.plan.analyze.Analysis; +import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; +import org.apache.iotdb.db.mpp.plan.analyze.schema.SchemaValidator; import org.apache.iotdb.db.mpp.plan.statement.Statement; +import org.apache.iotdb.rpc.RpcUtils; +import org.apache.iotdb.rpc.TSStatusCode; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; +import org.apache.iotdb.tsfile.write.schema.MeasurementSchema; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.Arrays; import java.util.Collections; import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.stream.Collectors; public abstract class InsertBaseStatement extends Statement { + private static final Logger LOGGER = LoggerFactory.getLogger(InsertBaseStatement.class); + + private static final PerformanceOverviewMetrics PERFORMANCE_OVERVIEW_METRICS = + PerformanceOverviewMetrics.getInstance(); + /** * if use id table, this filed is id form of device path <br> * if not, this filed is device path<br> @@ -35,10 +60,15 @@ public abstract class InsertBaseStatement extends Statement { protected boolean isAligned; + protected MeasurementSchema[] measurementSchemas; + protected String[] measurements; // get from client protected TSDataType[] dataTypes; + /** index of failed measurements -> info including measurement, data type and value */ + protected Map<Integer, FailedMeasurementInfo> failedMeasurementIndex2Info; + public PartialPath getDevicePath() { return devicePath; } @@ -55,12 +85,12 @@ public abstract class InsertBaseStatement extends Statement { this.measurements = measurements; } - public TSDataType[] getDataTypes() { - return dataTypes; + public MeasurementSchema[] getMeasurementSchemas() { + return measurementSchemas; } - public void setDataTypes(TSDataType[] dataTypes) { - this.dataTypes = dataTypes; + public void setMeasurementSchemas(MeasurementSchema[] measurementSchemas) { + this.measurementSchemas = measurementSchemas; } public boolean isAligned() { @@ -71,6 +101,14 @@ public abstract class InsertBaseStatement extends Statement { isAligned = aligned; } + public TSDataType[] getDataTypes() { + return dataTypes; + } + + public void setDataTypes(TSDataType[] dataTypes) { + this.dataTypes = dataTypes; + } + /** Returns true when this statement is empty and no need to write into the server */ public abstract boolean isEmpty(); @@ -78,4 +116,184 @@ public abstract class InsertBaseStatement extends Statement { public List<PartialPath> getPaths() { return Collections.emptyList(); } + + public abstract ISchemaValidation getSchemaValidation(); + + public abstract List<ISchemaValidation> getSchemaValidationList(); + + public void updateAfterSchemaValidation() throws QueryProcessException {} + + public void validateSchema(Analysis analysis) { + final long startTime = System.nanoTime(); + try { + SchemaValidator.validate(this); + } catch (SemanticException e) { + analysis.setFinishQueryAfterAnalyze(true); + if (e.getCause() instanceof IoTDBException) { + IoTDBException ioTDBException = (IoTDBException) e.getCause(); + analysis.setFailStatus( + RpcUtils.getStatus(ioTDBException.getErrorCode(), ioTDBException.getMessage())); + } else { + analysis.setFailStatus(RpcUtils.getStatus(TSStatusCode.METADATA_ERROR, e.getMessage())); + } + return; + } finally { + PERFORMANCE_OVERVIEW_METRICS.recordScheduleSchemaValidateCost(System.nanoTime() - startTime); + } + boolean hasFailedMeasurement = hasFailedMeasurements(); + String partialInsertMessage; + if (hasFailedMeasurement) { + partialInsertMessage = + String.format( + "Fail to insert measurements %s caused by %s", + getFailedMeasurements(), getFailedMessages()); + LOGGER.warn(partialInsertMessage); + analysis.setFailStatus( + RpcUtils.getStatus(TSStatusCode.METADATA_ERROR.getStatusCode(), partialInsertMessage)); + } + } + + /** Check whether data types are matched with measurement schemas */ + protected void selfCheckDataTypes(int index) + throws DataTypeMismatchException, PathNotExistException { + if (IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert()) { + // if enable partial insert, mark failed measurements with exception + if (measurementSchemas[index] == null) { + markFailedMeasurement( + index, + new PathNotExistException(devicePath.concatNode(measurements[index]).getFullPath())); + } else if ((dataTypes[index] != measurementSchemas[index].getType() + && !checkAndCastDataType(index, measurementSchemas[index].getType()))) { + markFailedMeasurement( + index, + new DataTypeMismatchException( + devicePath.getFullPath(), + measurements[index], + dataTypes[index], + measurementSchemas[index].getType(), + getMinTime(), + getFirstValueOfIndex(index))); + } + } else { + // if not enable partial insert, throw the exception directly + if (measurementSchemas[index] == null) { + throw new PathNotExistException(devicePath.concatNode(measurements[index]).getFullPath()); + } else if ((dataTypes[index] != measurementSchemas[index].getType() + && !checkAndCastDataType(index, measurementSchemas[index].getType()))) { + throw new DataTypeMismatchException( + devicePath.getFullPath(), + measurements[index], + dataTypes[index], + measurementSchemas[index].getType(), + getMinTime(), + getFirstValueOfIndex(index)); + } + } + } + + protected abstract boolean checkAndCastDataType(int columnIndex, TSDataType dataType); + + public abstract long getMinTime(); + + public abstract Object getFirstValueOfIndex(int index); + + // region partial insert + /** + * Mark failed measurement, measurements[index], dataTypes[index] and values/columns[index] would + * be null. We'd better use "measurements[index] == null" to determine if the measurement failed. + * <br> + * This method is not concurrency-safe. + * + * @param index failed measurement index + * @param cause cause Exception of failure + */ + public void markFailedMeasurement(int index, Exception cause) { + throw new UnsupportedOperationException(); + } + + public boolean hasValidMeasurements() { + for (Object o : measurements) { + if (o != null) { + return true; + } + } + return false; + } + + public boolean hasFailedMeasurements() { + return failedMeasurementIndex2Info != null && !failedMeasurementIndex2Info.isEmpty(); + } + + public int getFailedMeasurementNumber() { + return failedMeasurementIndex2Info == null ? 0 : failedMeasurementIndex2Info.size(); + } + + public List<String> getFailedMeasurements() { + return failedMeasurementIndex2Info == null + ? Collections.emptyList() + : failedMeasurementIndex2Info.values().stream() + .map(info -> info.measurement) + .collect(Collectors.toList()); + } + + public List<Exception> getFailedExceptions() { + return failedMeasurementIndex2Info == null + ? Collections.emptyList() + : failedMeasurementIndex2Info.values().stream() + .map(info -> info.cause) + .collect(Collectors.toList()); + } + + public List<String> getFailedMessages() { + return failedMeasurementIndex2Info == null + ? Collections.emptyList() + : failedMeasurementIndex2Info.values().stream() + .map( + info -> { + Throwable cause = info.cause; + while (cause.getCause() != null) { + cause = cause.getCause(); + } + return cause.getMessage(); + }) + .collect(Collectors.toList()); + } + + protected static class FailedMeasurementInfo { + protected String measurement; + protected TSDataType dataType; + protected Object value; + protected Exception cause; + + public FailedMeasurementInfo( + String measurement, TSDataType dataType, Object value, Exception cause) { + this.measurement = measurement; + this.dataType = dataType; + this.value = value; + this.cause = cause; + } + } + // endregion + + @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; + InsertBaseStatement that = (InsertBaseStatement) o; + return isAligned == that.isAligned + && Objects.equals(devicePath, that.devicePath) + && Arrays.equals(measurementSchemas, that.measurementSchemas) + && Arrays.equals(measurements, that.measurements) + && Arrays.equals(dataTypes, that.dataTypes); + } + + @Override + public int hashCode() { + int result = Objects.hash(super.hashCode(), devicePath, isAligned); + result = 31 * result + Arrays.hashCode(measurementSchemas); + result = 31 * result + Arrays.hashCode(measurements); + result = 31 * result + Arrays.hashCode(dataTypes); + return result; + } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertMultiTabletsStatement.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertMultiTabletsStatement.java index 77090231e6c..2ee89dc69c4 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertMultiTabletsStatement.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertMultiTabletsStatement.java @@ -20,12 +20,15 @@ package org.apache.iotdb.db.mpp.plan.statement.crud; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; import org.apache.iotdb.db.mpp.plan.statement.StatementType; import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor; +import org.apache.iotdb.tsfile.exception.NotImplementedException; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import java.util.ArrayList; import java.util.List; +import java.util.stream.Collectors; public class InsertMultiTabletsStatement extends InsertBaseStatement { @@ -95,4 +98,31 @@ public class InsertMultiTabletsStatement extends InsertBaseStatement { } return result; } + + @Override + public ISchemaValidation getSchemaValidation() { + throw new UnsupportedOperationException(); + } + + @Override + public List<ISchemaValidation> getSchemaValidationList() { + return insertTabletStatementList.stream() + .map(InsertTabletStatement::getSchemaValidation) + .collect(Collectors.toList()); + } + + @Override + protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { + return false; + } + + @Override + public long getMinTime() { + throw new NotImplementedException(); + } + + @Override + public Object getFirstValueOfIndex(int index) { + throw new NotImplementedException(); + } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowStatement.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowStatement.java index 7fc5659d90e..4015151d25c 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowStatement.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowStatement.java @@ -20,19 +20,38 @@ package org.apache.iotdb.db.mpp.plan.statement.crud; import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; +import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.db.conf.IoTDBDescriptor; +import org.apache.iotdb.db.exception.metadata.AlignedTimeseriesException; +import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException; +import org.apache.iotdb.db.exception.metadata.PathNotExistException; import org.apache.iotdb.db.exception.query.QueryProcessException; +import org.apache.iotdb.db.exception.sql.SemanticException; +import org.apache.iotdb.db.mpp.common.schematree.IMeasurementSchemaInfo; +import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; import org.apache.iotdb.db.mpp.plan.statement.StatementType; import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor; +import org.apache.iotdb.db.utils.CommonUtils; import org.apache.iotdb.db.utils.TimePartitionUtils; +import org.apache.iotdb.db.utils.TypeInferenceUtils; +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.apache.iotdb.tsfile.utils.ReadWriteIOUtils; +import org.apache.iotdb.tsfile.write.schema.MeasurementSchema; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; -public class InsertRowStatement extends InsertBaseStatement { +public class InsertRowStatement extends InsertBaseStatement implements ISchemaValidation { + + private static final Logger LOGGER = LoggerFactory.getLogger(InsertRowStatement.class); private static final byte TYPE_RAW_STRING = -1; private static final byte TYPE_NULL = -2; @@ -130,4 +149,165 @@ public class InsertRowStatement extends InsertBaseStatement { public <R, C> R accept(StatementVisitor<R, C> visitor, C context) { return visitor.visitInsertRow(this, context); } + + @Override + public long getMinTime() { + return getTime(); + } + + @Override + public Object getFirstValueOfIndex(int index) { + return values[index]; + } + + @Override + protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { + if (CommonUtils.checkCanCastType(dataTypes[columnIndex], dataType)) { + LOGGER.warn( + "Inserting to {}.{} : Cast from {} to {}", + devicePath, + measurements[columnIndex], + dataTypes[columnIndex], + dataType); + values[columnIndex] = + CommonUtils.castValue(dataTypes[columnIndex], dataType, values[columnIndex]); + dataTypes[columnIndex] = dataType; + return true; + } + return false; + } + + /** + * transfer String[] values to specific data types when isNeedInferType is true. <br> + * Notice: measurementSchemas must be initialized before calling this method + */ + @SuppressWarnings("squid:S3776") // Suppress high Cognitive Complexity warning + public void transferType() throws QueryProcessException { + + for (int i = 0; i < measurementSchemas.length; i++) { + // null when time series doesn't exist + if (measurementSchemas[i] == null) { + if (!IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert()) { + throw new QueryProcessException( + new PathNotExistException( + devicePath.getFullPath() + IoTDBConstant.PATH_SEPARATOR + measurements[i])); + } else { + markFailedMeasurement( + i, + new QueryProcessException( + new PathNotExistException( + devicePath.getFullPath() + IoTDBConstant.PATH_SEPARATOR + measurements[i]))); + } + continue; + } + // parse string value to specific type + dataTypes[i] = measurementSchemas[i].getType(); + try { + values[i] = CommonUtils.parseValue(dataTypes[i], values[i].toString()); + } catch (Exception e) { + LOGGER.warn( + "data type of {}.{} is not consistent, registered type {}, inserting timestamp {}, value {}", + devicePath, + measurements[i], + dataTypes[i], + time, + values[i]); + if (!IoTDBDescriptor.getInstance().getConfig().isEnablePartialInsert()) { + throw e; + } else { + markFailedMeasurement(i, e); + } + } + } + isNeedInferType = false; + } + + @Override + public void markFailedMeasurement(int index, Exception cause) { + if (measurements[index] == null) { + return; + } + + if (failedMeasurementIndex2Info == null) { + failedMeasurementIndex2Info = new HashMap<>(); + } + + InsertBaseStatement.FailedMeasurementInfo failedMeasurementInfo = + new InsertBaseStatement.FailedMeasurementInfo( + measurements[index], dataTypes[index], values[index], cause); + failedMeasurementIndex2Info.putIfAbsent(index, failedMeasurementInfo); + + measurements[index] = null; + dataTypes[index] = null; + values[index] = null; + } + + @Override + public ISchemaValidation getSchemaValidation() { + return this; + } + + @Override + public List<ISchemaValidation> getSchemaValidationList() { + throw new UnsupportedOperationException(); + } + + @Override + public void updateAfterSchemaValidation() throws QueryProcessException { + if (isNeedInferType) { + transferType(); + } + } + + @Override + public TSDataType getDataType(int index) { + if (isNeedInferType) { + return TypeInferenceUtils.getPredictedDataType(values[index], true); + } else { + return dataTypes[index]; + } + } + + @Override + public TSEncoding getEncoding(int index) { + return null; + } + + @Override + public CompressionType getCompressionType(int index) { + return null; + } + + @Override + public void validateDeviceSchema(boolean isAligned) { + if (this.isAligned != isAligned) { + throw new SemanticException( + new AlignedTimeseriesException( + String.format( + "timeseries under this device are%s aligned, " + "please use %s interface", + isAligned ? "" : " not", isAligned ? "aligned" : "non-aligned"), + devicePath.getFullPath())); + } + } + + @Override + public void validateMeasurementSchema(int index, IMeasurementSchemaInfo measurementSchemaInfo) { + if (measurementSchemas == null) { + measurementSchemas = new MeasurementSchema[measurements.length]; + } + if (measurementSchemaInfo == null) { + measurementSchemas[index] = null; + } else { + measurementSchemas[index] = measurementSchemaInfo.getSchema(); + } + if (isNeedInferType) { + return; + } + + try { + selfCheckDataTypes(index); + } catch (DataTypeMismatchException | PathNotExistException e) { + throw new SemanticException(e); + } + } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsOfOneDeviceStatement.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsOfOneDeviceStatement.java index 1507dd112b2..65c1f106a90 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsOfOneDeviceStatement.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsOfOneDeviceStatement.java @@ -21,14 +21,19 @@ package org.apache.iotdb.db.mpp.plan.statement.crud; import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.db.exception.query.QueryProcessException; +import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; import org.apache.iotdb.db.mpp.plan.statement.StatementType; import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor; import org.apache.iotdb.db.utils.TimePartitionUtils; +import org.apache.iotdb.tsfile.exception.NotImplementedException; +import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; +import java.util.stream.Collectors; public class InsertRowsOfOneDeviceStatement extends InsertBaseStatement { @@ -94,4 +99,41 @@ public class InsertRowsOfOneDeviceStatement extends InsertBaseStatement { } return ret; } + + @Override + public ISchemaValidation getSchemaValidation() { + throw new UnsupportedOperationException(); + } + + @Override + public List<ISchemaValidation> getSchemaValidationList() { + return insertRowStatementList.stream() + .map(InsertRowStatement::getSchemaValidation) + .collect(Collectors.toList()); + } + + @Override + public void updateAfterSchemaValidation() throws QueryProcessException { + for (InsertRowStatement insertRowStatement : insertRowStatementList) { + insertRowStatement.updateAfterSchemaValidation(); + if (!this.hasFailedMeasurements() && insertRowStatement.hasFailedMeasurements()) { + this.failedMeasurementIndex2Info = insertRowStatement.failedMeasurementIndex2Info; + } + } + } + + @Override + protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { + return false; + } + + @Override + public long getMinTime() { + throw new NotImplementedException(); + } + + @Override + public Object getFirstValueOfIndex(int index) { + throw new NotImplementedException(); + } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsStatement.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsStatement.java index c2b24e69261..0f0f973b947 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsStatement.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertRowsStatement.java @@ -20,12 +20,16 @@ package org.apache.iotdb.db.mpp.plan.statement.crud; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.db.exception.query.QueryProcessException; +import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; import org.apache.iotdb.db.mpp.plan.statement.StatementType; import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor; +import org.apache.iotdb.tsfile.exception.NotImplementedException; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; import java.util.ArrayList; import java.util.List; +import java.util.stream.Collectors; public class InsertRowsStatement extends InsertBaseStatement { @@ -95,4 +99,41 @@ public class InsertRowsStatement extends InsertBaseStatement { } return result; } + + @Override + public ISchemaValidation getSchemaValidation() { + throw new UnsupportedOperationException(); + } + + @Override + public List<ISchemaValidation> getSchemaValidationList() { + return insertRowStatementList.stream() + .map(InsertRowStatement::getSchemaValidation) + .collect(Collectors.toList()); + } + + @Override + public void updateAfterSchemaValidation() throws QueryProcessException { + for (InsertRowStatement insertRowStatement : insertRowStatementList) { + insertRowStatement.updateAfterSchemaValidation(); + if (!this.hasFailedMeasurements() && insertRowStatement.hasFailedMeasurements()) { + this.failedMeasurementIndex2Info = insertRowStatement.failedMeasurementIndex2Info; + } + } + } + + @Override + protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { + return false; + } + + @Override + public long getMinTime() { + throw new NotImplementedException(); + } + + @Override + public Object getFirstValueOfIndex(int index) { + throw new NotImplementedException(); + } } diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertTabletStatement.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertTabletStatement.java index 2da234041a0..5306c7939a7 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertTabletStatement.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/crud/InsertTabletStatement.java @@ -20,15 +20,36 @@ package org.apache.iotdb.db.mpp.plan.statement.crud; import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; import org.apache.iotdb.commons.path.PartialPath; +import org.apache.iotdb.db.exception.metadata.AlignedTimeseriesException; +import org.apache.iotdb.db.exception.metadata.DataTypeMismatchException; +import org.apache.iotdb.db.exception.metadata.PathNotExistException; +import org.apache.iotdb.db.exception.sql.SemanticException; +import org.apache.iotdb.db.mpp.common.schematree.IMeasurementSchemaInfo; +import org.apache.iotdb.db.mpp.plan.analyze.schema.ISchemaValidation; import org.apache.iotdb.db.mpp.plan.statement.StatementType; import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor; +import org.apache.iotdb.db.utils.CommonUtils; import org.apache.iotdb.db.utils.TimePartitionUtils; +import org.apache.iotdb.tsfile.exception.write.UnSupportedDataTypeException; +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.apache.iotdb.tsfile.utils.Binary; import org.apache.iotdb.tsfile.utils.BitMap; +import org.apache.iotdb.tsfile.write.schema.MeasurementSchema; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; -public class InsertTabletStatement extends InsertBaseStatement { +public class InsertTabletStatement extends InsertBaseStatement implements ISchemaValidation { + + private static final Logger LOGGER = LoggerFactory.getLogger(InsertTabletStatement.class); + + private static final String DATATYPE_UNSUPPORTED = "Data type %s is not supported."; private long[] times; // times should be sorted. It is done in the session API. private BitMap[] bitMaps; @@ -117,4 +138,136 @@ public class InsertTabletStatement extends InsertBaseStatement { } return ret; } + + @Override + public ISchemaValidation getSchemaValidation() { + return this; + } + + @Override + public List<ISchemaValidation> getSchemaValidationList() { + throw new UnsupportedOperationException(); + } + + @Override + protected boolean checkAndCastDataType(int columnIndex, TSDataType dataType) { + if (CommonUtils.checkCanCastType(dataTypes[columnIndex], dataType)) { + LOGGER.warn( + "Inserting to {}.{} : Cast from {} to {}", + devicePath, + measurements[columnIndex], + dataTypes[columnIndex], + dataType); + columns[columnIndex] = + CommonUtils.castArray(dataTypes[columnIndex], dataType, columns[columnIndex]); + dataTypes[columnIndex] = dataType; + return true; + } + return false; + } + + @Override + public void markFailedMeasurement(int index, Exception cause) { + if (measurements[index] == null) { + return; + } + + if (failedMeasurementIndex2Info == null) { + failedMeasurementIndex2Info = new HashMap<>(); + } + + InsertBaseStatement.FailedMeasurementInfo failedMeasurementInfo = + new InsertBaseStatement.FailedMeasurementInfo( + measurements[index], dataTypes[index], columns[index], cause); + failedMeasurementIndex2Info.putIfAbsent(index, failedMeasurementInfo); + + measurements[index] = null; + dataTypes[index] = null; + columns[index] = null; + } + + @Override + public long getMinTime() { + return times[0]; + } + + @Override + public Object getFirstValueOfIndex(int index) { + Object value; + switch (dataTypes[index]) { + case INT32: + int[] intValues = (int[]) columns[index]; + value = intValues[0]; + break; + case INT64: + long[] longValues = (long[]) columns[index]; + value = longValues[0]; + break; + case FLOAT: + float[] floatValues = (float[]) columns[index]; + value = floatValues[0]; + break; + case DOUBLE: + double[] doubleValues = (double[]) columns[index]; + value = doubleValues[0]; + break; + case BOOLEAN: + boolean[] boolValues = (boolean[]) columns[index]; + value = boolValues[0]; + break; + case TEXT: + Binary[] binaryValues = (Binary[]) columns[index]; + value = binaryValues[0]; + break; + default: + throw new UnSupportedDataTypeException( + String.format(DATATYPE_UNSUPPORTED, dataTypes[index])); + } + return value; + } + + @Override + public TSDataType getDataType(int index) { + return null; + } + + @Override + public TSEncoding getEncoding(int index) { + return null; + } + + @Override + public CompressionType getCompressionType(int index) { + return null; + } + + @Override + public void validateDeviceSchema(boolean isAligned) { + if (this.isAligned != isAligned) { + throw new SemanticException( + new AlignedTimeseriesException( + String.format( + "timeseries under this device are%s aligned, " + "please use %s interface", + isAligned ? "" : " not", isAligned ? "aligned" : "non-aligned"), + devicePath.getFullPath())); + } + } + + @Override + public void validateMeasurementSchema(int index, IMeasurementSchemaInfo measurementSchemaInfo) { + if (measurementSchemas == null) { + measurementSchemas = new MeasurementSchema[measurements.length]; + } + if (measurementSchemaInfo == null) { + measurementSchemas[index] = null; + } else { + measurementSchemas[index] = measurementSchemaInfo.getSchema(); + } + + try { + selfCheckDataTypes(index); + } catch (DataTypeMismatchException | PathNotExistException e) { + throw new SemanticException(e); + } + } } diff --git a/server/src/main/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoer.java b/server/src/main/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoer.java index 89430bc91f2..093aa30e2aa 100644 --- a/server/src/main/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoer.java +++ b/server/src/main/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoer.java @@ -27,7 +27,6 @@ import org.apache.iotdb.db.engine.storagegroup.TsFileResource; import org.apache.iotdb.db.exception.WriteProcessException; import org.apache.iotdb.db.metadata.idtable.IDTable; import org.apache.iotdb.db.metadata.idtable.entry.DeviceIDFactory; -import org.apache.iotdb.db.mpp.plan.analyze.schema.SchemaValidator; import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.DeleteDataNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertNode; import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowNode; @@ -109,9 +108,6 @@ public class TsFilePlanRedoer { // TODO get device id by idTable // idTable.getSeriesSchemas(node); } else { - if (!IoTDBDescriptor.getInstance().getConfig().isClusterMode()) { - SchemaValidator.validate(node); - } node.setDeviceID(DeviceIDFactory.getInstance().getDeviceID(node.getDevicePath())); } diff --git a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java index 34b2655e6cc..89406f8f580 100644 --- a/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java +++ b/server/src/test/java/org/apache/iotdb/db/engine/storagegroup/DataRegionTest.java @@ -282,11 +282,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode1.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode1); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -304,11 +304,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode2.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode2); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -445,11 +445,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode1.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode1); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -466,11 +466,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode2.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode2); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -533,11 +533,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode1.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode1); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -554,11 +554,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode2.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode2); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -621,11 +621,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode1.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode1); dataRegion.asyncCloseAllWorkingTsFileProcessors(); @@ -642,11 +642,11 @@ public class DataRegionTest { false, measurements, dataTypes, + measurementSchemas, times, null, columns, times.length); - insertTabletNode2.setMeasurementSchemas(measurementSchemas); dataRegion.insertTablet(insertTabletNode2); dataRegion.asyncCloseAllWorkingTsFileProcessors(); diff --git a/server/src/test/java/org/apache/iotdb/db/wal/io/WALFileTest.java b/server/src/test/java/org/apache/iotdb/db/wal/io/WALFileTest.java index bb696fdda84..1156021d7f5 100644 --- a/server/src/test/java/org/apache/iotdb/db/wal/io/WALFileTest.java +++ b/server/src/test/java/org/apache/iotdb/db/wal/io/WALFileTest.java @@ -236,18 +236,6 @@ public class WALFileTest { } bitMaps[i].mark(i % times.length); } - - InsertTabletNode insertTabletNode = - new InsertTabletNode( - new PlanNodeId(""), - new PartialPath(devicePath), - false, - new String[] {"s1", "s2", "s3", "s4", "s5", "s6"}, - dataTypes, - times, - bitMaps, - columns, - times.length); MeasurementSchema[] schemas = new MeasurementSchema[] { new MeasurementSchema("s1", dataTypes[0]), @@ -257,9 +245,18 @@ public class WALFileTest { new MeasurementSchema("s5", dataTypes[4]), new MeasurementSchema("s6", dataTypes[5]), }; - insertTabletNode.setMeasurementSchemas(schemas); - return insertTabletNode; + return new InsertTabletNode( + new PlanNodeId(""), + new PartialPath(devicePath), + false, + new String[] {"s1", "s2", "s3", "s4", "s5", "s6"}, + dataTypes, + schemas, + times, + bitMaps, + columns, + times.length); } public static DeleteDataNode getDeleteDataNode(String devicePath) throws IllegalPathException { diff --git a/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java b/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java index 517b7aeaf1c..508338a6e07 100644 --- a/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java +++ b/server/src/test/java/org/apache/iotdb/db/wal/node/ConsensusReqReaderTest.java @@ -697,19 +697,12 @@ public class ConsensusReqReaderTest { bitMaps[i].mark(i % times.length); } - InsertTabletNode insertTabletNode = - new InsertTabletNode( - new PlanNodeId(""), - new PartialPath(devicePath), - false, - new String[] {"s1", "s2", "s3", "s4", "s5", "s6"}, - dataTypes, - times, - bitMaps, - columns, - times.length); - - insertTabletNode.setMeasurementSchemas( + return new InsertTabletNode( + new PlanNodeId(""), + new PartialPath(devicePath), + false, + new String[] {"s1", "s2", "s3", "s4", "s5", "s6"}, + dataTypes, new MeasurementSchema[] { new MeasurementSchema("s1", TSDataType.DOUBLE), new MeasurementSchema("s2", TSDataType.FLOAT), @@ -717,9 +710,11 @@ public class ConsensusReqReaderTest { new MeasurementSchema("s4", TSDataType.INT32), new MeasurementSchema("s5", TSDataType.BOOLEAN), new MeasurementSchema("s6", TSDataType.TEXT) - }); - - return insertTabletNode; + }, + times, + bitMaps, + columns, + times.length); } private DeleteDataNode getDeleteDataNode(String devicePath) throws IllegalPathException { diff --git a/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java b/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java index 8542c817ef7..48523d4c0e0 100644 --- a/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java +++ b/server/src/test/java/org/apache/iotdb/db/wal/node/WALNodeTest.java @@ -196,24 +196,22 @@ public class WALNodeTest { } bitMaps[i].mark(i % times.length); } - - InsertTabletNode insertTabletNode = - new InsertTabletNode( - new PlanNodeId(""), - new PartialPath(devicePath), - false, - measurements, - dataTypes, - times, - bitMaps, - columns, - times.length); MeasurementSchema[] schemas = new MeasurementSchema[6]; for (int i = 0; i < 6; i++) { schemas[i] = new MeasurementSchema(measurements[i], dataTypes[i], TSEncoding.PLAIN); } - insertTabletNode.setMeasurementSchemas(schemas); - return insertTabletNode; + + return new InsertTabletNode( + new PlanNodeId(""), + new PartialPath(devicePath), + false, + measurements, + dataTypes, + schemas, + times, + bitMaps, + columns, + times.length); } @Test diff --git a/server/src/test/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoerTest.java b/server/src/test/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoerTest.java index 91843aecf6b..dffbf321815 100644 --- a/server/src/test/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/wal/recover/file/TsFilePlanRedoerTest.java @@ -291,6 +291,7 @@ public class TsFilePlanRedoerTest { false, new String[] {"s1", "s2"}, new TSDataType[] {TSDataType.INT32, TSDataType.INT64}, + null, times, bitMaps, columns, @@ -392,6 +393,7 @@ public class TsFilePlanRedoerTest { TSDataType.FLOAT, TSDataType.TEXT }, + null, times, bitMaps, columns, @@ -473,6 +475,10 @@ public class TsFilePlanRedoerTest { false, new String[] {"s1", "s2"}, new TSDataType[] {TSDataType.INT32, TSDataType.INT64}, + new MeasurementSchema[] { + new MeasurementSchema("s1", TSDataType.INT32), + new MeasurementSchema("s2", TSDataType.INT64), + }, times, null, columns, @@ -520,15 +526,14 @@ public class TsFilePlanRedoerTest { false, new String[] {"s1", "s2"}, new TSDataType[] {TSDataType.INT32, TSDataType.INT64}, + new MeasurementSchema[] { + new MeasurementSchema("s1", TSDataType.INT32), + new MeasurementSchema("s2", TSDataType.INT64), + }, times, null, columns, times.length); - insertTabletNode.setMeasurementSchemas( - new MeasurementSchema[] { - new MeasurementSchema("s1", TSDataType.INT32), - new MeasurementSchema("s2", TSDataType.INT64), - }); // redo InsertTabletPlan, vsg processor is used to test IdTable, don't test IdTable here TsFilePlanRedoer planRedoer = new TsFilePlanRedoer(tsFileResource, false, null); @@ -634,6 +639,14 @@ public class TsFilePlanRedoerTest { // mark value of time=9 as null bitMaps[i].mark(3); } + MeasurementSchema[] schemas = + new MeasurementSchema[] { + null, + new MeasurementSchema("s2", TSDataType.INT64), + new MeasurementSchema("s3", TSDataType.BOOLEAN), + new MeasurementSchema("s4", TSDataType.FLOAT), + null + }; InsertTabletNode insertTabletNode = new InsertTabletNode( new PlanNodeId(""), @@ -647,20 +660,13 @@ public class TsFilePlanRedoerTest { TSDataType.FLOAT, TSDataType.TEXT }, + schemas, times, bitMaps, columns, times.length); // redo InsertTabletPlan, data region is used to test IdTable, don't test IdTable here TsFilePlanRedoer planRedoer = new TsFilePlanRedoer(tsFileResource, true, null); - MeasurementSchema[] schemas = - new MeasurementSchema[] { - null, - new MeasurementSchema("s2", TSDataType.INT64), - new MeasurementSchema("s3", TSDataType.BOOLEAN), - new MeasurementSchema("s4", TSDataType.FLOAT), - null - }; insertTabletNode.setMeasurementSchemas(schemas); planRedoer.redoInsert(insertTabletNode); diff --git a/server/src/test/java/org/apache/iotdb/db/wal/recover/file/UnsealedTsFileRecoverPerformerTest.java b/server/src/test/java/org/apache/iotdb/db/wal/recover/file/UnsealedTsFileRecoverPerformerTest.java index d1fe54b54b3..16d00c67dc3 100644 --- a/server/src/test/java/org/apache/iotdb/db/wal/recover/file/UnsealedTsFileRecoverPerformerTest.java +++ b/server/src/test/java/org/apache/iotdb/db/wal/recover/file/UnsealedTsFileRecoverPerformerTest.java @@ -297,6 +297,7 @@ public class UnsealedTsFileRecoverPerformerTest { false, new String[] {"s1"}, new TSDataType[] {TSDataType.INT64}, + null, new long[] {time}, null, new Integer[] {1},
