This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch insertValidate in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 5c0f9293f022db66bffda68209e33a507bc7a89b Author: HTHou <[email protected]> AuthorDate: Tue Jun 14 15:48:51 2022 +0800 init --- .../iotdb/db/engine/storagegroup/DataRegion.java | 24 ---------------------- .../scheduler/FragmentInstanceDispatcherImpl.java | 13 ++++++++++-- .../service/thrift/impl/InternalServiceImpl.java | 13 ++++++++++-- 3 files changed, 22 insertions(+), 28 deletions(-) diff --git a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java index bd9f80b0c9..2d5a4a8e16 100755 --- a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/DataRegion.java @@ -874,14 +874,6 @@ public class DataRegion { } finally { writeUnlock(); } - - if (insertRowNode.hasFailedMeasurements()) { - logger.warn( - "Fail to insert measurements {} caused by {}", - insertRowNode.getFailedMeasurements(), - insertRowNode.getFailedMessages()); - checkFailedMeasurements(insertRowNode); - } } /** @@ -1096,14 +1088,6 @@ public class DataRegion { } finally { writeUnlock(); } - - if (insertTabletNode.hasFailedMeasurements()) { - logger.warn( - "Fail to insert measurements {} caused by {}", - insertTabletNode.getFailedMeasurements(), - insertTabletNode.getFailedMessages()); - checkFailedMeasurements(insertTabletNode); - } } /** @@ -3500,14 +3484,6 @@ public class DataRegion { } } - private void checkFailedMeasurements(InsertNode node) throws WriteProcessException { - List<Exception> exceptions = node.getFailedExceptions(); - throw new WriteProcessException( - "failed to insert measurements " - + node.getFailedMeasurements() - + (!exceptions.isEmpty() ? (" caused by " + exceptions.get(0).getMessage()) : "")); - } - @TestOnly public long getPartitionMaxFileVersions(long partitionId) { return partitionMaxFileVersions.getOrDefault(partitionId, 0L); 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 c30d8f61ec..4ac3e83aa8 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 @@ -201,12 +201,21 @@ public class FragmentInstanceDispatcherImpl implements IFragInstanceDispatcher { return !((FragmentInstanceInfo) readResponse.getDataset()).getState().isFailed(); case WRITE: PlanNode planNode = instance.getFragment().getRoot(); + boolean hasFailedMeasurement = false; if (planNode instanceof InsertNode) { + InsertNode insertNode = (InsertNode) planNode; try { - SchemaValidator.validate((InsertNode) planNode); + SchemaValidator.validate(insertNode); } catch (SemanticException e) { throw new FragmentInstanceDispatchException(e); } + hasFailedMeasurement = insertNode.hasFailedMeasurements(); + if (hasFailedMeasurement) { + logger.warn( + "Fail to insert measurements {} caused by {}", + insertNode.getFailedMeasurements(), + insertNode.getFailedMessages()); + } } ConsensusWriteResponse writeResponse; if (groupId instanceof DataRegionId) { @@ -214,7 +223,7 @@ public class FragmentInstanceDispatcherImpl implements IFragInstanceDispatcher { } else { writeResponse = SchemaRegionConsensusImpl.getInstance().write(groupId, planNode); } - return TSStatusCode.SUCCESS_STATUS.getStatusCode() == writeResponse.getStatus().getCode(); + return !hasFailedMeasurement && TSStatusCode.SUCCESS_STATUS.getStatusCode() == writeResponse.getStatus().getCode(); } throw new UnsupportedOperationException( String.format("unknown query type [%s]", instance.getType())); diff --git a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java index 0e9a00ecce..7c5999addc 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java @@ -144,14 +144,23 @@ public class InternalServiceImpl implements InternalService.Iface { ConsensusWriteResponse writeResponse; PlanNode planNode = PlanNodeType.deserialize(req.planNode.body); + boolean hasFailedMeasurement = false; if (planNode instanceof InsertNode) { + InsertNode insertNode = (InsertNode) planNode; try { - SchemaValidator.validate((InsertNode) planNode); + SchemaValidator.validate(insertNode); } catch (SemanticException e) { response.setAccepted(false); response.setMessage(e.getMessage()); return response; } + hasFailedMeasurement = insertNode.hasFailedMeasurements(); + if (hasFailedMeasurement) { + LOGGER.warn( + "Fail to insert measurements {} caused by {}", + insertNode.getFailedMeasurements(), + insertNode.getFailedMessages()); + } } if (groupId instanceof DataRegionId) { writeResponse = DataRegionConsensusImpl.getInstance().write(groupId, planNode); @@ -160,7 +169,7 @@ public class InternalServiceImpl implements InternalService.Iface { } // TODO need consider more status if (writeResponse.getStatus() != null) { - response.setAccepted( + response.setAccepted(!hasFailedMeasurement && TSStatusCode.SUCCESS_STATUS.getStatusCode() == writeResponse.getStatus().getCode()); response.setMessage(writeResponse.getStatus().message); } else {
