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

haonan pushed a commit to branch rel/0.12
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/rel/0.12 by this push:
     new b974543  [IOTDB-1632] fill only if the value is missing (#3914) (#3956)
b974543 is described below

commit b974543c530e533cddcd457787bec6ea6e06a544
Author: Zhong Wang <[email protected]>
AuthorDate: Tue Sep 14 13:22:27 2021 +0800

    [IOTDB-1632] fill only if the value is missing (#3914) (#3956)
---
 .../iotdb/cluster/query/ClusterQueryRouter.java    | 14 +---
 .../cluster/query/fill/ClusterFillExecutor.java    | 48 +++++++++--
 .../cluster/query/ClusterFillExecutorTest.java     | 22 ++---
 .../iotdb/db/query/executor/FillQueryExecutor.java | 94 +++++++++++++++++-----
 .../iotdb/db/query/executor/QueryRouter.java       | 21 +----
 .../apache/iotdb/db/integration/IoTDBFillIT.java   | 34 +++++++-
 6 files changed, 162 insertions(+), 71 deletions(-)

diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterQueryRouter.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterQueryRouter.java
index 71669aa..e3be92c 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterQueryRouter.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterQueryRouter.java
@@ -27,8 +27,8 @@ import 
org.apache.iotdb.cluster.query.last.ClusterLastQueryExecutor;
 import org.apache.iotdb.cluster.server.member.MetaGroupMember;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
-import org.apache.iotdb.db.metadata.PartialPath;
 import org.apache.iotdb.db.qp.physical.crud.AggregationPlan;
+import org.apache.iotdb.db.qp.physical.crud.FillQueryPlan;
 import org.apache.iotdb.db.qp.physical.crud.GroupByTimePlan;
 import org.apache.iotdb.db.qp.physical.crud.LastQueryPlan;
 import org.apache.iotdb.db.qp.physical.crud.RawDataQueryPlan;
@@ -41,9 +41,7 @@ import org.apache.iotdb.db.query.executor.FillQueryExecutor;
 import org.apache.iotdb.db.query.executor.LastQueryExecutor;
 import org.apache.iotdb.db.query.executor.QueryRouter;
 import org.apache.iotdb.db.query.executor.RawDataQueryExecutor;
-import org.apache.iotdb.db.query.executor.fill.IFill;
 import 
org.apache.iotdb.tsfile.exception.filter.QueryFilterOptimizationException;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 import org.apache.iotdb.tsfile.read.expression.ExpressionType;
 import org.apache.iotdb.tsfile.read.expression.IExpression;
 import org.apache.iotdb.tsfile.read.expression.util.ExpressionOptimizer;
@@ -51,8 +49,6 @@ import 
org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
 
 import java.io.IOException;
 import java.util.ArrayList;
-import java.util.List;
-import java.util.Map;
 
 public class ClusterQueryRouter extends QueryRouter {
 
@@ -63,12 +59,8 @@ public class ClusterQueryRouter extends QueryRouter {
   }
 
   @Override
-  protected FillQueryExecutor getFillExecutor(
-      List<PartialPath> fillPaths,
-      List<TSDataType> dataTypes,
-      long queryTime,
-      Map<TSDataType, IFill> fillType) {
-    return new ClusterFillExecutor(fillPaths, dataTypes, queryTime, fillType, 
metaGroupMember);
+  protected FillQueryExecutor getFillExecutor(FillQueryPlan plan) {
+    return new ClusterFillExecutor(plan, metaGroupMember);
   }
 
   @Override
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterFillExecutor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterFillExecutor.java
index f6cd707..1dc1363 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterFillExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterFillExecutor.java
@@ -19,31 +19,36 @@
 
 package org.apache.iotdb.cluster.query.fill;
 
+import org.apache.iotdb.cluster.query.reader.ClusterReaderFactory;
 import org.apache.iotdb.cluster.server.member.MetaGroupMember;
+import org.apache.iotdb.db.exception.StorageEngineException;
+import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.metadata.PartialPath;
+import org.apache.iotdb.db.qp.physical.crud.FillQueryPlan;
 import org.apache.iotdb.db.query.context.QueryContext;
 import org.apache.iotdb.db.query.executor.FillQueryExecutor;
 import org.apache.iotdb.db.query.executor.fill.IFill;
 import org.apache.iotdb.db.query.executor.fill.LinearFill;
 import org.apache.iotdb.db.query.executor.fill.PreviousFill;
+import org.apache.iotdb.db.query.reader.series.IReaderByTimestamp;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.read.TimeValuePair;
+import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
 
+import java.io.IOException;
+import java.util.ArrayList;
 import java.util.List;
-import java.util.Map;
 import java.util.Set;
 
 public class ClusterFillExecutor extends FillQueryExecutor {
 
   private MetaGroupMember metaGroupMember;
+  private ClusterReaderFactory clusterReaderFactory;
 
-  public ClusterFillExecutor(
-      List<PartialPath> selectedSeries,
-      List<TSDataType> dataTypes,
-      long queryTime,
-      Map<TSDataType, IFill> typeIFillMap,
-      MetaGroupMember metaGroupMember) {
-    super(selectedSeries, dataTypes, queryTime, typeIFillMap);
+  public ClusterFillExecutor(FillQueryPlan plan, MetaGroupMember 
metaGroupMember) {
+    super(plan);
     this.metaGroupMember = metaGroupMember;
+    this.clusterReaderFactory = new ClusterReaderFactory(metaGroupMember);
   }
 
   @Override
@@ -65,4 +70,31 @@ public class ClusterFillExecutor extends FillQueryExecutor {
     }
     return null;
   }
+
+  @Override
+  protected List<TimeValuePair> getTimeValuePairs(QueryContext context)
+      throws QueryProcessException, StorageEngineException, IOException {
+    List<TimeValuePair> ret = new ArrayList<>(selectedSeries.size());
+
+    for (int i = 0; i < selectedSeries.size(); i++) {
+      PartialPath path = selectedSeries.get(i);
+      TSDataType dataType = dataTypes.get(i);
+      IReaderByTimestamp reader =
+          clusterReaderFactory.getReaderByTimestamp(
+              path,
+              plan.getAllMeasurementsInDevice(path.getDevice()),
+              dataTypes.get(i),
+              context,
+              plan.isAscending());
+
+      Object[] results = reader.getValuesInTimestamps(new long[] {queryTime}, 
1);
+      if (results[0] != null) {
+        ret.add(new TimeValuePair(queryTime, 
TsPrimitiveType.getByType(dataType, results[0])));
+      } else {
+        ret.add(null);
+      }
+    }
+
+    return ret;
+  }
 }
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/query/ClusterFillExecutorTest.java
 
b/cluster/src/test/java/org/apache/iotdb/cluster/query/ClusterFillExecutorTest.java
index 8bbd9fa..5c8c176 100644
--- 
a/cluster/src/test/java/org/apache/iotdb/cluster/query/ClusterFillExecutorTest.java
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/query/ClusterFillExecutorTest.java
@@ -75,14 +75,9 @@ public class ClusterFillExecutorTest extends BaseQueryTest {
             new Object[] {10.0},
           };
       for (int i = 0; i < queryTimes.length; i++) {
-        fillExecutor =
-            new ClusterFillExecutor(
-                plan.getDeduplicatedPaths(),
-                plan.getDeduplicatedDataTypes(),
-                queryTimes[i],
-                plan.getFillType(),
-                testMetaMember);
-        queryDataSet = fillExecutor.execute(context, plan);
+        plan.setQueryTime(queryTimes[i]);
+        fillExecutor = new ClusterFillExecutor(plan, testMetaMember);
+        queryDataSet = fillExecutor.execute(context);
         checkDoubleDataset(queryDataSet, answers[i]);
         assertFalse(queryDataSet.hasNext());
       }
@@ -122,14 +117,9 @@ public class ClusterFillExecutorTest extends BaseQueryTest 
{
             new Object[] {null},
           };
       for (int i = 0; i < queryTimes.length; i++) {
-        fillExecutor =
-            new ClusterFillExecutor(
-                plan.getDeduplicatedPaths(),
-                plan.getDeduplicatedDataTypes(),
-                queryTimes[i],
-                plan.getFillType(),
-                testMetaMember);
-        queryDataSet = fillExecutor.execute(context, plan);
+        plan.setQueryTime(queryTimes[i]);
+        fillExecutor = new ClusterFillExecutor(plan, testMetaMember);
+        queryDataSet = fillExecutor.execute(context);
         checkDoubleDataset(queryDataSet, answers[i]);
         assertFalse(queryDataSet.hasNext());
       }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/executor/FillQueryExecutor.java
 
b/server/src/main/java/org/apache/iotdb/db/query/executor/FillQueryExecutor.java
index af97632..f78e1ae 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/query/executor/FillQueryExecutor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/query/executor/FillQueryExecutor.java
@@ -21,43 +21,49 @@ package org.apache.iotdb.db.query.executor;
 
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.engine.StorageEngine;
+import org.apache.iotdb.db.engine.querycontext.QueryDataSource;
 import org.apache.iotdb.db.engine.storagegroup.StorageGroupProcessor;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.metadata.PartialPath;
 import org.apache.iotdb.db.qp.physical.crud.FillQueryPlan;
 import org.apache.iotdb.db.query.context.QueryContext;
+import org.apache.iotdb.db.query.control.QueryResourceManager;
 import org.apache.iotdb.db.query.dataset.SingleDataSet;
 import org.apache.iotdb.db.query.executor.fill.IFill;
 import org.apache.iotdb.db.query.executor.fill.PreviousFill;
+import org.apache.iotdb.db.query.reader.series.ManagedSeriesReader;
+import org.apache.iotdb.db.query.reader.series.SeriesRawDataBatchReader;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 import org.apache.iotdb.tsfile.read.TimeValuePair;
+import org.apache.iotdb.tsfile.read.common.BatchData;
 import org.apache.iotdb.tsfile.read.common.RowRecord;
+import org.apache.iotdb.tsfile.read.filter.TimeFilter;
+import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 import org.apache.iotdb.tsfile.read.query.dataset.QueryDataSet;
 
 import javax.activation.UnsupportedDataTypeException;
 
 import java.io.IOException;
+import java.util.ArrayList;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
 
 public class FillQueryExecutor {
 
-  private List<PartialPath> selectedSeries;
-  private List<TSDataType> dataTypes;
-  private long queryTime;
-  private Map<TSDataType, IFill> typeIFillMap;
+  protected FillQueryPlan plan;
+  protected List<PartialPath> selectedSeries;
+  protected List<TSDataType> dataTypes;
+  protected Map<TSDataType, IFill> typeIFillMap;
+  protected long queryTime;
 
-  public FillQueryExecutor(
-      List<PartialPath> selectedSeries,
-      List<TSDataType> dataTypes,
-      long queryTime,
-      Map<TSDataType, IFill> typeIFillMap) {
-    this.selectedSeries = selectedSeries;
-    this.queryTime = queryTime;
-    this.typeIFillMap = typeIFillMap;
-    this.dataTypes = dataTypes;
+  public FillQueryExecutor(FillQueryPlan fillQueryPlan) {
+    this.plan = fillQueryPlan;
+    this.selectedSeries = plan.getDeduplicatedPaths();
+    this.typeIFillMap = plan.getFillType();
+    this.dataTypes = plan.getDeduplicatedDataTypes();
+    this.queryTime = plan.getQueryTime();
   }
 
   /**
@@ -65,18 +71,25 @@ public class FillQueryExecutor {
    *
    * @param context query context
    */
-  public QueryDataSet execute(QueryContext context, FillQueryPlan 
fillQueryPlan)
+  public QueryDataSet execute(QueryContext context)
       throws StorageEngineException, QueryProcessException, IOException {
     RowRecord record = new RowRecord(queryTime);
 
     List<StorageGroupProcessor> list = 
StorageEngine.getInstance().mergeLock(selectedSeries);
     try {
+      List<TimeValuePair> timeValuePairs = getTimeValuePairs(context);
+      long defaultFillInterval = 
IoTDBDescriptor.getInstance().getConfig().getDefaultFillInterval();
       for (int i = 0; i < selectedSeries.size(); i++) {
         PartialPath path = selectedSeries.get(i);
         TSDataType dataType = dataTypes.get(i);
+
+        if (timeValuePairs.get(i) != null) {
+          // No need to fill
+          record.addField(timeValuePairs.get(i).getValue().getValue(), 
dataType);
+          continue;
+        }
+
         IFill fill;
-        long defaultFillInterval =
-            IoTDBDescriptor.getInstance().getConfig().getDefaultFillInterval();
         if (!typeIFillMap.containsKey(dataType)) {
           switch (dataType) {
             case INT32:
@@ -88,7 +101,7 @@ public class FillQueryExecutor {
               fill = new PreviousFill(dataType, queryTime, 
defaultFillInterval);
               break;
             default:
-              throw new UnsupportedDataTypeException("do not support datatype 
" + dataType);
+              throw new UnsupportedDataTypeException("unsupported data type " 
+ dataType);
           }
         } else {
           fill = typeIFillMap.get(dataType).copy();
@@ -99,7 +112,7 @@ public class FillQueryExecutor {
                 path,
                 dataType,
                 queryTime,
-                fillQueryPlan.getAllMeasurementsInDevice(path.getDevice()),
+                plan.getAllMeasurementsInDevice(path.getDevice()),
                 context);
 
         TimeValuePair timeValuePair = fill.getFillResult();
@@ -128,4 +141,49 @@ public class FillQueryExecutor {
     fill.configureFill(path, dataType, queryTime, deviceMeasurements, context);
     return fill;
   }
+
+  protected List<TimeValuePair> getTimeValuePairs(QueryContext context)
+      throws QueryProcessException, StorageEngineException, IOException {
+    List<ManagedSeriesReader> readers = initManagedSeriesReader(context);
+    List<TimeValuePair> ret = new ArrayList<>(selectedSeries.size());
+    for (ManagedSeriesReader reader : readers) {
+      if (reader.hasNextBatch()) {
+        BatchData batchData = reader.nextBatch();
+        if (batchData.hasCurrent()) {
+          ret.add(new TimeValuePair(batchData.currentTime(), 
batchData.currentTsPrimitiveType()));
+          continue;
+        }
+      }
+      ret.add(null);
+    }
+
+    return ret;
+  }
+
+  private List<ManagedSeriesReader> initManagedSeriesReader(QueryContext 
context)
+      throws StorageEngineException, QueryProcessException {
+    Filter timeFilter = TimeFilter.eq(queryTime);
+    List<ManagedSeriesReader> readers = new ArrayList<>();
+    for (int i = 0; i < selectedSeries.size(); i++) {
+      PartialPath path = selectedSeries.get(i);
+      TSDataType dataType = dataTypes.get(i);
+      QueryDataSource queryDataSource =
+          QueryResourceManager.getInstance().getQueryDataSource(path, context, 
timeFilter);
+      timeFilter = queryDataSource.updateFilterUsingTTL(timeFilter);
+      ManagedSeriesReader reader =
+          new SeriesRawDataBatchReader(
+              path,
+              plan.getAllMeasurementsInDevice(path.getDevice()),
+              dataType,
+              context,
+              queryDataSource,
+              timeFilter,
+              null,
+              null,
+              plan.isAscending());
+      readers.add(reader);
+    }
+
+    return readers;
+  }
 }
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/executor/QueryRouter.java 
b/server/src/main/java/org/apache/iotdb/db/query/executor/QueryRouter.java
index 77bfd59..cb9da3d 100644
--- a/server/src/main/java/org/apache/iotdb/db/query/executor/QueryRouter.java
+++ b/server/src/main/java/org/apache/iotdb/db/query/executor/QueryRouter.java
@@ -36,9 +36,7 @@ import 
org.apache.iotdb.db.query.dataset.groupby.GroupByFillDataSet;
 import org.apache.iotdb.db.query.dataset.groupby.GroupByTimeDataSet;
 import org.apache.iotdb.db.query.dataset.groupby.GroupByWithValueFilterDataSet;
 import 
org.apache.iotdb.db.query.dataset.groupby.GroupByWithoutValueFilterDataSet;
-import org.apache.iotdb.db.query.executor.fill.IFill;
 import 
org.apache.iotdb.tsfile.exception.filter.QueryFilterOptimizationException;
-import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 import org.apache.iotdb.tsfile.read.expression.ExpressionType;
 import org.apache.iotdb.tsfile.read.expression.IExpression;
 import org.apache.iotdb.tsfile.read.expression.impl.BinaryExpression;
@@ -54,7 +52,6 @@ import org.slf4j.LoggerFactory;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
-import java.util.Map;
 
 /**
  * Query entrance class of IoTDB query process. All query clause will be 
transformed to physical
@@ -229,22 +226,12 @@ public class QueryRouter implements IQueryRouter {
   @Override
   public QueryDataSet fill(FillQueryPlan fillQueryPlan, QueryContext context)
       throws StorageEngineException, QueryProcessException, IOException {
-    List<PartialPath> fillPaths = fillQueryPlan.getDeduplicatedPaths();
-    List<TSDataType> dataTypes = fillQueryPlan.getDeduplicatedDataTypes();
-    long queryTime = fillQueryPlan.getQueryTime();
-    Map<TSDataType, IFill> fillType = fillQueryPlan.getFillType();
-
-    FillQueryExecutor fillQueryExecutor =
-        getFillExecutor(fillPaths, dataTypes, queryTime, fillType);
-    return fillQueryExecutor.execute(context, fillQueryPlan);
+    FillQueryExecutor fillQueryExecutor = getFillExecutor(fillQueryPlan);
+    return fillQueryExecutor.execute(context);
   }
 
-  protected FillQueryExecutor getFillExecutor(
-      List<PartialPath> fillPaths,
-      List<TSDataType> dataTypes,
-      long queryTime,
-      Map<TSDataType, IFill> fillType) {
-    return new FillQueryExecutor(fillPaths, dataTypes, queryTime, fillType);
+  protected FillQueryExecutor getFillExecutor(FillQueryPlan plan) {
+    return new FillQueryExecutor(plan);
   }
 
   @Override
diff --git 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBFillIT.java 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBFillIT.java
index ec527dc..6ee10c5 100644
--- a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBFillIT.java
+++ b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBFillIT.java
@@ -312,7 +312,7 @@ public class IoTDBFillIT {
   }
 
   @Test
-  public void ValueFillTest() {
+  public void valueFillTest() {
     String res = "7,7.0,true,7";
     try (Connection connection =
             DriverManager.getConnection("jdbc:iotdb://127.0.0.1:6667/", 
"root", "root");
@@ -344,6 +344,38 @@ public class IoTDBFillIT {
   }
 
   @Test
+  public void valueFillNonNullTest() {
+    String res = "1,1.1,false,11";
+    try (Connection connection =
+            DriverManager.getConnection("jdbc:iotdb://127.0.0.1:6667/", 
"root", "root");
+        Statement statement = connection.createStatement()) {
+
+      boolean hasResultSet =
+          statement.execute(
+              "SELECT temperature, status, hardware"
+                  + " FROM root.ln.wf01.wt01"
+                  + " WHERE time = 1 FILL(int32[7], double[7], 
boolean[true])");
+
+      Assert.assertTrue(hasResultSet);
+      ResultSet resultSet = statement.getResultSet();
+      while (resultSet.next()) {
+        String ans =
+            resultSet.getString(TIMESTAMP_STR)
+                + ","
+                + resultSet.getString(TEMPERATURE_STR_1)
+                + ","
+                + resultSet.getString(STATUS_STR_1)
+                + ","
+                + resultSet.getString(HARDWARE_STR);
+        Assert.assertEquals(res, ans);
+      }
+    } catch (Exception e) {
+      e.printStackTrace();
+      fail(e.getMessage());
+    }
+  }
+
+  @Test
   public void PreviousFillTest() {
     String[] retArray1 = new String[] {"3,3.3,false,33", "70,50.5,false,550", 
"70,null,null,null"};
     try (Connection connection =

Reply via email to