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 =
