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(