fix bug in session part
diff --git a/example/session/src/main/java/org/apache/iotdb/FastInsertExample.java b/example/session/src/main/java/org/apache/iotdb/FastInsertExample.java index 5e2a3a9..b24dd72 100644 --- a/example/session/src/main/java/org/apache/iotdb/FastInsertExample.java +++ b/example/session/src/main/java/org/apache/iotdb/FastInsertExample.java
@@ -65,7 +65,7 @@ List<Long> timestamps = new ArrayList<>(); List<List<TSDataType>> typesList = new ArrayList<>(); - for (long time = 0; time < 500; time++) { + for (long time = 1000; time < 1500; time++) { List<Object> values = new ArrayList<>(); List<TSDataType> types = new ArrayList<>(); values.add(1L);
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 f876697..c4f7077 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
@@ -257,8 +257,6 @@ LOGGER.warn(partialInsertMessage); } - // TODO: (FAStWRITE) (侯昊男) 根据数据类型把 values 反序列化出来 - // TODO: (FAStWRITE) 然后再进入到这一步 ConsensusWriteResponse writeResponse = fireTriggerAndInsert(context.getRegionId(), insertNode);
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 2d6b10e..b9270d3 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
@@ -43,6 +43,7 @@ if (insertNode instanceof FastInsertRowsNode) { SCHEMA_FETCHER.fetchAndComputeSchemaWithAutoCreateForFastWrite( ((BatchInsertNode) insertNode).getSchemaValidationList()); + ((FastInsertRowsNode) insertNode).fillValues(); } else { SCHEMA_FETCHER.fetchAndComputeSchemaWithAutoCreate( ((BatchInsertNode) insertNode).getSchemaValidationList());
diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowNode.java index 80ea5de..378d6a2 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowNode.java
@@ -111,6 +111,10 @@ this.rawValues = ByteBuffer.wrap(bytes); } + public boolean hasFailedMeasurements() { + return false; + } + @Override public void initMeasurementSchemaContainer(int size, String[] measurements) { this.measurementSchemas = new MeasurementSchema[size]; @@ -131,4 +135,32 @@ measurementSchemas[index] = measurementSchemaInfo.getSchema(); this.dataTypes[index] = measurementSchemaInfo.getSchema().getType(); } + + public void fillValues() throws QueryProcessException { + this.values = new Object[measurements.length]; + for (int i = 0; i < dataTypes.length; i++) { + switch (dataTypes[i]) { + case BOOLEAN: + values[i] = ReadWriteIOUtils.readBool(this.rawValues); + break; + case INT32: + values[i] = ReadWriteIOUtils.readInt(this.rawValues); + break; + case INT64: + values[i] = ReadWriteIOUtils.readLong(this.rawValues); + break; + case FLOAT: + values[i] = ReadWriteIOUtils.readFloat(this.rawValues); + break; + case DOUBLE: + values[i] = ReadWriteIOUtils.readDouble(this.rawValues); + break; + case TEXT: + values[i] = ReadWriteIOUtils.readBinary(this.rawValues); + break; + default: + throw new QueryProcessException("Unsupported data type:" + dataTypes[i]); + } + } + } }
diff --git a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowsNode.java b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowsNode.java index aaa3815..2106660 100644 --- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowsNode.java +++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/write/FastInsertRowsNode.java
@@ -20,6 +20,7 @@ package org.apache.iotdb.db.mpp.plan.planner.plan.node.write; import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; +import org.apache.iotdb.db.exception.query.QueryProcessException; import org.apache.iotdb.db.mpp.plan.analyze.Analysis; import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.mpp.plan.planner.plan.node.WritePlanNode; @@ -42,6 +43,16 @@ super(id, insertRowNodeIndexList, fastInsertRowNodeList); } + public boolean hasFailedMeasurements() { + return false; + } + + public void fillValues() throws QueryProcessException { + for (int i = 0; i < getInsertRowNodeList().size(); i++) { + ((FastInsertRowNode) getInsertRowNodeList().get(i)).fillValues(); + } + } + @Override public List<WritePlanNode> splitByPartition(Analysis analysis) { Map<TRegionReplicaSet, InsertRowsNode> splitMap = new HashMap<>();
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 74a35f5..d4d883e 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
@@ -77,7 +77,7 @@ private long time; // TODO: (FASTWRITE) (侯昊男) 增加 byteBuffer 字段 - private Object[] values; + protected Object[] values; private boolean isNeedInferType = false;
diff --git a/session/src/main/java/org/apache/iotdb/session/Session.java b/session/src/main/java/org/apache/iotdb/session/Session.java index e39b236..822ba5c 100644 --- a/session/src/main/java/org/apache/iotdb/session/Session.java +++ b/session/src/main/java/org/apache/iotdb/session/Session.java
@@ -2480,8 +2480,7 @@ throws IoTDBConnectionException { request.addToPrefixPaths(deviceId); request.addToTimestamps(time); - ByteBuffer buffer = - ByteBuffer.allocate(SessionUtils.calculateLengthForFastInsert(types, values)); + ByteBuffer buffer = SessionUtils.getValueBuffer(types, values); request.addToValuesList(buffer); }
diff --git a/session/src/main/java/org/apache/iotdb/session/util/SessionUtils.java b/session/src/main/java/org/apache/iotdb/session/util/SessionUtils.java index 5f55031b..8002773 100644 --- a/session/src/main/java/org/apache/iotdb/session/util/SessionUtils.java +++ b/session/src/main/java/org/apache/iotdb/session/util/SessionUtils.java
@@ -79,6 +79,7 @@ public static ByteBuffer getValueBuffer(List<TSDataType> types, List<Object> values) throws IoTDBConnectionException { ByteBuffer buffer = ByteBuffer.allocate(SessionUtils.calculateLength(types, values)); + SessionUtils.putValues(types, values, buffer); return buffer; }