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 =