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

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

commit 61999aa1555337c6c486f484ad5bd08318508358
Author: 151250176 <[email protected]>
AuthorDate: Mon Dec 13 11:02:14 2021 +0800

    fix normal query
---
 .../iotdb/cluster/query/LocalQueryExecutor.java    |  8 +++++-
 .../query/last/ClusterLastQueryExecutor.java       |  8 +++++-
 .../engine/storagegroup/StorageGroupProcessor.java |  1 +
 .../db/engine/storagegroup/TsFileProcessor.java    |  8 +++---
 .../id_table/AppendOnlyDiskSchemaManager.java      | 19 ++++++++++++++
 .../apache/iotdb/db/metadata/id_table/IDTable.java | 28 ++++++++++++++++++++
 .../db/metadata/id_table/entry/SHA256DeviceID.java | 30 +++++++++++++++++++++-
 .../apache/iotdb/db/qp/executor/PlanExecutor.java  | 17 ++++++++++++
 .../iotdb/db/query/reader/series/SeriesReader.java |  5 ++--
 .../db/metadata/id_table/QueryWithIDTableTest.java | 30 +++++++++++-----------
 10 files changed, 130 insertions(+), 24 deletions(-)

diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/LocalQueryExecutor.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/LocalQueryExecutor.java
index dd6ec9f..dee1563 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/LocalQueryExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/LocalQueryExecutor.java
@@ -41,6 +41,7 @@ import org.apache.iotdb.cluster.rpc.thrift.RaftNode;
 import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
 import org.apache.iotdb.cluster.server.member.DataGroupMember;
 import org.apache.iotdb.cluster.utils.ClusterUtils;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.metadata.IllegalPathException;
 import org.apache.iotdb.db.exception.metadata.MetadataException;
@@ -1027,7 +1028,12 @@ public class LocalQueryExecutor {
 
     List<Pair<Boolean, TimeValuePair>> timeValuePairs =
         LastQueryExecutor.calculateLastPairForSeriesLocally(
-            partialPaths, dataTypes, queryContext, expression, 
request.getDeviceMeasurements());
+            partialPaths,
+            dataTypes,
+            queryContext,
+            expression,
+            request.getDeviceMeasurements(),
+            IoTDBDescriptor.getInstance().getConfig().isEnableIDTable());
     ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
     DataOutputStream dataOutputStream = new 
DataOutputStream(byteArrayOutputStream);
     for (Pair<Boolean, TimeValuePair> timeValuePair : timeValuePairs) {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
index db6b14d..4debef1 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
@@ -32,6 +32,7 @@ import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.server.member.DataGroupMember;
 import org.apache.iotdb.cluster.server.member.MetaGroupMember;
 import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.metadata.path.PartialPath;
@@ -198,7 +199,12 @@ public class ClusterLastQueryExecutor extends 
LastQueryExecutor {
         throw new QueryProcessException(e.getMessage());
       }
       return calculateLastPairForSeriesLocally(
-          seriesPaths, dataTypes, context, expression, 
queryPlan.getDeviceToMeasurements());
+          seriesPaths,
+          dataTypes,
+          context,
+          expression,
+          queryPlan.getDeviceToMeasurements(),
+          IoTDBDescriptor.getInstance().getConfig().isEnableIDTable());
     }
 
     private List<Pair<Boolean, TimeValuePair>> calculateSeriesLastRemotely(
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
index ce88214..bd454f5 100755
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
@@ -1742,6 +1742,7 @@ public class StorageGroupProcessor {
       Filter timeFilter)
       throws QueryProcessException {
     readLock();
+    fullPath = IDTable.translateQueryPath(fullPath);
     try {
       List<TsFileResource> seqResources =
           getFileResourceListForQuery(
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
index 9b18306..bc04d6c 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
@@ -245,12 +245,12 @@ public class TsFileProcessor {
 
     // update start time of this memtable
     tsFileResource.updateStartTime(
-        insertRowPlan.getDeviceId().getFullPath(), insertRowPlan.getTime());
+        insertRowPlan.getDeviceID().toStringID(), insertRowPlan.getTime());
     // for sequence tsfile, we update the endTime only when the file is 
prepared to be closed.
     // for unsequence tsfile, we have to update the endTime for each insertion.
     if (!sequence) {
       tsFileResource.updateEndTime(
-          insertRowPlan.getDeviceId().getFullPath(), insertRowPlan.getTime());
+          insertRowPlan.getDeviceID().toStringID(), insertRowPlan.getTime());
     }
     tsFileResource.updatePlanIndexes(insertRowPlan.getIndex());
   }
@@ -327,13 +327,13 @@ public class TsFileProcessor {
       results[i] = RpcUtils.SUCCESS_STATUS;
     }
     tsFileResource.updateStartTime(
-        insertTabletPlan.getDeviceId().getFullPath(), 
insertTabletPlan.getTimes()[start]);
+        insertTabletPlan.getDeviceID().toStringID(), 
insertTabletPlan.getTimes()[start]);
 
     // for sequence tsfile, we update the endTime only when the file is 
prepared to be closed.
     // for unsequence tsfile, we have to update the endTime for each insertion.
     if (!sequence) {
       tsFileResource.updateEndTime(
-          insertTabletPlan.getDeviceId().getFullPath(), 
insertTabletPlan.getTimes()[end - 1]);
+          insertTabletPlan.getDeviceID().toStringID(), 
insertTabletPlan.getTimes()[end - 1]);
     }
     tsFileResource.updatePlanIndexes(insertTabletPlan.getIndex());
   }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/id_table/AppendOnlyDiskSchemaManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/id_table/AppendOnlyDiskSchemaManager.java
index 5ba6ef0..77d32aa 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/id_table/AppendOnlyDiskSchemaManager.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/id_table/AppendOnlyDiskSchemaManager.java
@@ -1,3 +1,22 @@
+/*
+ * 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.metadata.id_table;
 
 import org.apache.iotdb.db.metadata.id_table.entry.DiskSchemaEntry;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/id_table/IDTable.java 
b/server/src/main/java/org/apache/iotdb/db/metadata/id_table/IDTable.java
index 7b336d2..1ef1586 100644
--- a/server/src/main/java/org/apache/iotdb/db/metadata/id_table/IDTable.java
+++ b/server/src/main/java/org/apache/iotdb/db/metadata/id_table/IDTable.java
@@ -30,6 +30,7 @@ import 
org.apache.iotdb.db.metadata.id_table.entry.InsertMeasurementMNode;
 import org.apache.iotdb.db.metadata.id_table.entry.SchemaEntry;
 import org.apache.iotdb.db.metadata.id_table.entry.TimeseriesID;
 import org.apache.iotdb.db.metadata.mnode.IMeasurementMNode;
+import org.apache.iotdb.db.metadata.path.MeasurementPath;
 import org.apache.iotdb.db.metadata.path.PartialPath;
 import org.apache.iotdb.db.qp.physical.crud.InsertPlan;
 import org.apache.iotdb.db.qp.physical.sys.CreateAlignedTimeSeriesPlan;
@@ -177,6 +178,8 @@ public class IDTable {
 
     // set reusable device id
     plan.setDeviceID(deviceEntry.getDeviceID());
+    // for last flushed time map
+    plan.setDeviceId(new PartialPath(deviceEntry.getDeviceID().toStringID()));
 
     return deviceEntry.getDeviceID();
   }
@@ -409,6 +412,31 @@ public class IDTable {
     }
   }
 
+  /**
+   * translate query path's device path to device id
+   *
+   * @param fullPath full query path
+   * @return translated query path
+   */
+  public static PartialPath translateQueryPath(PartialPath fullPath) {
+    // if not enable id table, just return original path
+    if (!config.isEnableIDTable()) {
+      return fullPath;
+    }
+
+    TimeseriesID timeseriesID = new TimeseriesID(fullPath);
+    try {
+      return new MeasurementPath(
+          timeseriesID.getDeviceID().toStringID(),
+          timeseriesID.getMeasurement(),
+          fullPath.getMeasurementSchema());
+    } catch (MetadataException e) {
+      logger.error("Error when translate query path: " + fullPath);
+    }
+
+    return null;
+  }
+
   @TestOnly
   public Map<IDeviceID, DeviceEntry>[] getIdTables() {
     return idTables;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/id_table/entry/SHA256DeviceID.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/id_table/entry/SHA256DeviceID.java
index 2fc42bf..0e7b30b 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/id_table/entry/SHA256DeviceID.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/id_table/entry/SHA256DeviceID.java
@@ -28,6 +28,8 @@ public class SHA256DeviceID implements IDeviceID {
   long l3;
   long l4;
 
+  private static final String SEPERATOR = "#";
+
   private static MessageDigest md;
 
   static {
@@ -39,6 +41,32 @@ public class SHA256DeviceID implements IDeviceID {
   }
 
   public SHA256DeviceID(String deviceID) {
+    if (deviceID.indexOf('.') == -1) {
+      fromSHA256String(deviceID);
+    } else {
+      buildSHA256(deviceID);
+    }
+  }
+
+  /**
+   * build device id from a sha 256 string, like "1#1#1#1"
+   *
+   * @param deviceID a sha 256 string
+   */
+  private void fromSHA256String(String deviceID) {
+    String[] part = deviceID.split(SEPERATOR);
+    l1 = Long.parseLong(part[0]);
+    l2 = Long.parseLong(part[1]);
+    l3 = Long.parseLong(part[2]);
+    l4 = Long.parseLong(part[3]);
+  }
+
+  /**
+   * build device id from a device path
+   *
+   * @param deviceID device path
+   */
+  private void buildSHA256(String deviceID) {
     byte[] hashVal = md.digest(deviceID.getBytes());
     md.reset();
 
@@ -82,6 +110,6 @@ public class SHA256DeviceID implements IDeviceID {
 
   @Override
   public String toStringID() {
-    return l1 + "#" + l2 + "#" + l3 + "#" + l4;
+    return l1 + SEPERATOR + l2 + SEPERATOR + l3 + SEPERATOR + l4;
   }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java 
b/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
index b726262..76e804e 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
@@ -607,6 +607,23 @@ public class PlanExecutor implements IPlanExecutor {
       } else if (queryPlan instanceof LastQueryPlan) {
         queryDataSet = queryRouter.lastQuery((LastQueryPlan) queryPlan, 
context);
       } else {
+        RawDataQueryPlan plan = (RawDataQueryPlan) queryPlan;
+        //        List<PartialPath> list = new ArrayList<>();
+        //        for (PartialPath path : plan.getDeduplicatedPaths()) {
+        //          TimeseriesID timeseriesID = new TimeseriesID(path);
+        //          try {
+        //            PartialPath fullPath =
+        //                new PartialPath(
+        //                    timeseriesID.getDeviceID().toStringID(),
+        // timeseriesID.getMeasurement());
+        //            list.add(fullPath);
+        //          } catch (IllegalPathException e) {
+        //            e.printStackTrace();
+        //          }
+        //        }
+        //        plan.setDeduplicatedPaths(list);
+        //        plan.setDeduplicatedVectorPaths(list);
+
         queryDataSet = queryRouter.rawDataQuery((RawDataQueryPlan) queryPlan, 
context);
       }
     }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
 
b/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
index 9515552..c9a4ccd 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/reader/series/SeriesReader.java
@@ -20,6 +20,7 @@ package org.apache.iotdb.db.query.reader.series;
 
 import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
 import org.apache.iotdb.db.engine.storagegroup.TsFileResource;
+import org.apache.iotdb.db.metadata.id_table.IDTable;
 import org.apache.iotdb.db.metadata.path.PartialPath;
 import org.apache.iotdb.db.query.context.QueryContext;
 import org.apache.iotdb.db.query.control.QueryTimeManager;
@@ -149,7 +150,7 @@ public class SeriesReader {
       Filter valueFilter,
       TsFileFilter fileFilter,
       boolean ascending) {
-    this.seriesPath = seriesPath;
+    this.seriesPath = IDTable.translateQueryPath(seriesPath);
     this.allSensors = allSensors;
     this.dataType = dataType;
     this.context = context;
@@ -192,7 +193,7 @@ public class SeriesReader {
       Filter timeFilter,
       Filter valueFilter,
       boolean ascending) {
-    this.seriesPath = seriesPath;
+    this.seriesPath = IDTable.translateQueryPath(seriesPath);
     this.allSensors = allSensors;
     this.allSensors.add(seriesPath.getMeasurement());
     this.dataType = dataType;
diff --git 
a/server/src/test/java/org/apache/iotdb/db/metadata/id_table/QueryWithIDTableTest.java
 
b/server/src/test/java/org/apache/iotdb/db/metadata/id_table/QueryWithIDTableTest.java
index 7ea7579..39baabf 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/metadata/id_table/QueryWithIDTableTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/metadata/id_table/QueryWithIDTableTest.java
@@ -72,12 +72,12 @@ public class QueryWithIDTableTest {
   Set<String> retSet =
       new HashSet<>(
           Arrays.asList(
-              "113\troot.isp.d1.s3\t10003\tINT64",
-              "113\troot.isp.d1.s4\t103\tINT32",
+              "113\troot.isp.d1.s3\t100003\tINT64",
+              "113\troot.isp.d1.s4\t1003\tINT32",
               "113\troot.isp.d1.s5\tfalse\tBOOLEAN",
-              "113\troot.isp.d1.s6\thh3\tTEXT",
-              "113\troot.isp.d1.s1\t4.0\tDOUBLE",
-              "113\troot.isp.d1.s2\t5.0\tFLOAT"));
+              "113\troot.isp.d1.s6\tmm3\tTEXT",
+              "113\troot.isp.d1.s1\t13.0\tDOUBLE",
+              "113\troot.isp.d1.s2\t23.0\tFLOAT"));
 
   @Before
   public void before() {
@@ -142,25 +142,25 @@ public class QueryWithIDTableTest {
 
     // test it from id table
     assertEquals(
-        new TimeValuePair(113L, new TsDouble(4.0d)),
+        new TimeValuePair(113L, new TsDouble(13.0d)),
         StorageEngine.getInstance()
             .getProcessor(new PartialPath("root.isp.d1"))
             .getIdTable()
             .getLastCache(new TimeseriesID(new 
PartialPath("root.isp.d1.s1"))));
     assertEquals(
-        new TimeValuePair(113L, new TsFloat(5.0f)),
+        new TimeValuePair(113L, new TsFloat(23.0f)),
         StorageEngine.getInstance()
             .getProcessor(new PartialPath("root.isp.d1"))
             .getIdTable()
             .getLastCache(new TimeseriesID(new 
PartialPath("root.isp.d1.s2"))));
     assertEquals(
-        new TimeValuePair(113L, new TsLong(10003L)),
+        new TimeValuePair(113L, new TsLong(100003L)),
         StorageEngine.getInstance()
             .getProcessor(new PartialPath("root.isp.d1"))
             .getIdTable()
             .getLastCache(new TimeseriesID(new 
PartialPath("root.isp.d1.s3"))));
     assertEquals(
-        new TimeValuePair(113L, new TsInt(103)),
+        new TimeValuePair(113L, new TsInt(1003)),
         StorageEngine.getInstance()
             .getProcessor(new PartialPath("root.isp.d1"))
             .getIdTable()
@@ -172,7 +172,7 @@ public class QueryWithIDTableTest {
             .getIdTable()
             .getLastCache(new TimeseriesID(new 
PartialPath("root.isp.d1.s5"))));
     assertEquals(
-        new TimeValuePair(113L, new TsBinary(new Binary("hh3"))),
+        new TimeValuePair(113L, new TsBinary(new Binary("mm3"))),
         StorageEngine.getInstance()
             .getProcessor(new PartialPath("root.isp.d1"))
             .getIdTable()
@@ -198,12 +198,12 @@ public class QueryWithIDTableTest {
     columns[5] = new Binary[4];
 
     for (int r = 0; r < 4; r++) {
-      ((double[]) columns[0])[r] = 1.0 + r;
-      ((float[]) columns[1])[r] = 2 + r;
-      ((long[]) columns[2])[r] = 10000 + r;
-      ((int[]) columns[3])[r] = 100 + r;
+      ((double[]) columns[0])[r] = 10.0 + r;
+      ((float[]) columns[1])[r] = 20 + r;
+      ((long[]) columns[2])[r] = 100000 + r;
+      ((int[]) columns[3])[r] = 1000 + r;
       ((boolean[]) columns[4])[r] = false;
-      ((Binary[]) columns[5])[r] = new Binary("hh" + r);
+      ((Binary[]) columns[5])[r] = new Binary("mm" + r);
     }
 
     InsertTabletPlan tabletPlan =

Reply via email to