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

qiaojialin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new b3355c58b3 [IOTDB-3485] Insert with wrong type didn't return error 
massage (#6277)
b3355c58b3 is described below

commit b3355c58b3e74a154c623b5fbe38b621d47ed383
Author: Haonan <[email protected]>
AuthorDate: Tue Jun 14 19:22:26 2022 +0800

    [IOTDB-3485] Insert with wrong type didn't return error massage (#6277)
---
 .../db/it/aligned/IoTDBInsertAlignedValuesIT.java  |  9 ++------
 .../iotdb/db/engine/storagegroup/DataRegion.java   | 25 ----------------------
 .../scheduler/FragmentInstanceDispatcherImpl.java  | 14 ++++++++++--
 .../service/thrift/impl/InternalServiceImpl.java   | 15 +++++++++++--
 4 files changed, 27 insertions(+), 36 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/db/it/aligned/IoTDBInsertAlignedValuesIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/db/it/aligned/IoTDBInsertAlignedValuesIT.java
index 11d0d1b646..3840230cb5 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/db/it/aligned/IoTDBInsertAlignedValuesIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/db/it/aligned/IoTDBInsertAlignedValuesIT.java
@@ -240,19 +240,14 @@ public class IoTDBInsertAlignedValuesIT {
     }
   }
 
-  // TODO remove Ignore annotation while fixing this bug
-  @Ignore
-  @Test
-  public void testInsertWithWrongType() {
+  @Test(expected = Exception.class)
+  public void testInsertWithWrongType() throws SQLException {
     try (Connection connection = EnvFactory.getEnv().getConnection();
         Statement statement = connection.createStatement()) {
       statement.execute(
           "CREATE ALIGNED TIMESERIES root.lz.dev.GPS(latitude INT32 
encoding=PLAIN compressor=SNAPPY, longitude INT32 encoding=PLAIN 
compressor=SNAPPY) ");
       statement.execute(
           "insert into root.lz.dev.GPS(time,latitude,longitude) aligned 
values(1,1.3,6.7)");
-      fail();
-    } catch (SQLException e) {
-      assertEquals(313, e.getErrorCode());
     }
   }
 
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..4ce9ab0064 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
@@ -65,7 +65,6 @@ import org.apache.iotdb.db.metadata.idtable.IDTable;
 import org.apache.iotdb.db.metadata.idtable.IDTableManager;
 import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertMultiTabletsNode;
-import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowNode;
 import org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowsNode;
 import 
org.apache.iotdb.db.mpp.plan.planner.plan.node.write.InsertRowsOfOneDeviceNode;
@@ -874,14 +873,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 +1087,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 +3483,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..f994366f3f 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,8 @@ 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..17efadbf68 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);
@@ -161,7 +170,9 @@ public class InternalServiceImpl implements 
InternalService.Iface {
     // TODO need consider more status
     if (writeResponse.getStatus() != null) {
       response.setAccepted(
-          TSStatusCode.SUCCESS_STATUS.getStatusCode() == 
writeResponse.getStatus().getCode());
+          !hasFailedMeasurement
+              && TSStatusCode.SUCCESS_STATUS.getStatusCode()
+                  == writeResponse.getStatus().getCode());
       response.setMessage(writeResponse.getStatus().message);
     } else {
       LOGGER.error(

Reply via email to