This is an automated email from the ASF dual-hosted git repository.
jackietien pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new aaadc5c Fix exception when getting Statistics of aligned time series
in memory (#4430)
aaadc5c is described below
commit aaadc5ca5b0bc0bfc88b2d89a7818af8d90d271b
Author: liuminghui233 <[email protected]>
AuthorDate: Fri Nov 19 16:07:23 2021 +0800
Fix exception when getting Statistics of aligned time series in memory
(#4430)
---
.../querycontext/AlignedReadOnlyMemChunk.java | 4 +-
.../query/reader/chunk/MemAlignedChunkLoader.java | 52 ++++++++++
.../query/reader/chunk/MemAlignedChunkReader.java | 110 +++++++++++++++++++++
.../query/reader/chunk/MemAlignedPageReader.java | 91 +++++++++++++++++
.../iotdb/db/query/reader/series/SeriesReader.java | 16 +--
.../tsfile/read/reader/IAlignedPageReader.java | 26 +++++
.../tsfile/read/reader/page/AlignedPageReader.java | 4 +-
7 files changed, 292 insertions(+), 11 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
index ddc1cc7..2f18ef7 100644
---
a/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/querycontext/AlignedReadOnlyMemChunk.java
@@ -20,7 +20,7 @@
package org.apache.iotdb.db.engine.querycontext;
import org.apache.iotdb.db.exception.query.QueryProcessException;
-import org.apache.iotdb.db.query.reader.chunk.MemChunkLoader;
+import org.apache.iotdb.db.query.reader.chunk.MemAlignedChunkLoader;
import org.apache.iotdb.db.utils.datastructure.AlignedTVList;
import org.apache.iotdb.db.utils.datastructure.TVList;
import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
@@ -151,7 +151,7 @@ public class AlignedReadOnlyMemChunk extends
ReadOnlyMemChunk {
}
IChunkMetadata vectorChunkMetadata =
new AlignedChunkMetadata(timeChunkMetadata, valueChunkMetadataList);
- vectorChunkMetadata.setChunkLoader(new MemChunkLoader(this));
+ vectorChunkMetadata.setChunkLoader(new MemAlignedChunkLoader(this));
vectorChunkMetadata.setVersion(Long.MAX_VALUE);
cachedMetaData = vectorChunkMetadata;
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedChunkLoader.java
b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedChunkLoader.java
new file mode 100644
index 0000000..b307af1
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedChunkLoader.java
@@ -0,0 +1,52 @@
+/*
+ * 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.query.reader.chunk;
+
+import org.apache.iotdb.db.engine.querycontext.AlignedReadOnlyMemChunk;
+import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata;
+import org.apache.iotdb.tsfile.file.metadata.IChunkMetadata;
+import org.apache.iotdb.tsfile.read.common.Chunk;
+import org.apache.iotdb.tsfile.read.controller.IChunkLoader;
+import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+import org.apache.iotdb.tsfile.read.reader.IChunkReader;
+
+/** To read one aligned chunk from memory, and only used in iotdb server
module */
+public class MemAlignedChunkLoader implements IChunkLoader {
+
+ private final AlignedReadOnlyMemChunk chunk;
+
+ public MemAlignedChunkLoader(AlignedReadOnlyMemChunk chunk) {
+ this.chunk = chunk;
+ }
+
+ @Override
+ public Chunk loadChunk(ChunkMetadata chunkMetaData) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public void close() {
+ // no resources need to close
+ }
+
+ @Override
+ public IChunkReader getChunkReader(IChunkMetadata chunkMetaData, Filter
timeFilter) {
+ return new MemAlignedChunkReader(chunk, timeFilter);
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedChunkReader.java
b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedChunkReader.java
new file mode 100644
index 0000000..f52602b
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedChunkReader.java
@@ -0,0 +1,110 @@
+/*
+ * 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.query.reader.chunk;
+
+import org.apache.iotdb.db.engine.querycontext.AlignedReadOnlyMemChunk;
+import org.apache.iotdb.tsfile.file.metadata.AlignedChunkMetadata;
+import org.apache.iotdb.tsfile.read.TimeValuePair;
+import org.apache.iotdb.tsfile.read.common.BatchData;
+import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+import org.apache.iotdb.tsfile.read.reader.IChunkReader;
+import org.apache.iotdb.tsfile.read.reader.IPageReader;
+import org.apache.iotdb.tsfile.read.reader.IPointReader;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.List;
+
+/** To read aligned chunk data in memory */
+public class MemAlignedChunkReader implements IChunkReader, IPointReader {
+
+ private IPointReader timeValuePairIterator;
+ private Filter filter;
+ private boolean hasCachedTimeValuePair;
+ private TimeValuePair cachedTimeValuePair;
+ private List<IPageReader> pageReaderList;
+
+ public MemAlignedChunkReader(AlignedReadOnlyMemChunk readableChunk, Filter
filter) {
+ timeValuePairIterator = readableChunk.getPointReader();
+ this.filter = filter;
+ // we treat one ReadOnlyMemChunk as one Page
+ this.pageReaderList =
+ Collections.singletonList(
+ new MemAlignedPageReader(
+ timeValuePairIterator,
+ (AlignedChunkMetadata) readableChunk.getChunkMetaData(),
+ filter));
+ }
+
+ @Override
+ public boolean hasNextTimeValuePair() throws IOException {
+ if (hasCachedTimeValuePair) {
+ return true;
+ }
+ while (timeValuePairIterator.hasNextTimeValuePair()) {
+ TimeValuePair timeValuePair = timeValuePairIterator.nextTimeValuePair();
+ if (filter == null
+ || filter.satisfy(timeValuePair.getTimestamp(),
timeValuePair.getValue().getValue())) {
+ hasCachedTimeValuePair = true;
+ cachedTimeValuePair = timeValuePair;
+ break;
+ }
+ }
+ return hasCachedTimeValuePair;
+ }
+
+ @Override
+ public TimeValuePair nextTimeValuePair() throws IOException {
+ if (hasCachedTimeValuePair) {
+ hasCachedTimeValuePair = false;
+ return cachedTimeValuePair;
+ } else {
+ return timeValuePairIterator.nextTimeValuePair();
+ }
+ }
+
+ @Override
+ public TimeValuePair currentTimeValuePair() throws IOException {
+ if (!hasCachedTimeValuePair) {
+ cachedTimeValuePair = timeValuePairIterator.nextTimeValuePair();
+ hasCachedTimeValuePair = true;
+ }
+ return cachedTimeValuePair;
+ }
+
+ @Override
+ public boolean hasNextSatisfiedPage() throws IOException {
+ return hasNextTimeValuePair();
+ }
+
+ @Override
+ public BatchData nextPageData() throws IOException {
+ return pageReaderList.remove(0).getAllSatisfiedPageData();
+ }
+
+ @Override
+ public void close() {
+ // Do nothing because mem chunk reader will not open files
+ }
+
+ @Override
+ public List<IPageReader> loadPageReaderList() {
+ return this.pageReaderList;
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java
b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java
new file mode 100644
index 0000000..a084f12
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/query/reader/chunk/MemAlignedPageReader.java
@@ -0,0 +1,91 @@
+/*
+ * 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.query.reader.chunk;
+
+import org.apache.iotdb.tsfile.file.metadata.AlignedChunkMetadata;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics;
+import org.apache.iotdb.tsfile.read.TimeValuePair;
+import org.apache.iotdb.tsfile.read.common.BatchData;
+import org.apache.iotdb.tsfile.read.common.BatchDataFactory;
+import org.apache.iotdb.tsfile.read.filter.basic.Filter;
+import org.apache.iotdb.tsfile.read.filter.operator.AndFilter;
+import org.apache.iotdb.tsfile.read.reader.IAlignedPageReader;
+import org.apache.iotdb.tsfile.read.reader.IPageReader;
+import org.apache.iotdb.tsfile.read.reader.IPointReader;
+
+import java.io.IOException;
+
+public class MemAlignedPageReader implements IPageReader, IAlignedPageReader {
+
+ private final IPointReader timeValuePairIterator;
+ private final AlignedChunkMetadata chunkMetadata;
+ private Filter valueFilter;
+
+ public MemAlignedPageReader(
+ IPointReader timeValuePairIterator, AlignedChunkMetadata chunkMetadata,
Filter filter) {
+ this.timeValuePairIterator = timeValuePairIterator;
+ this.chunkMetadata = chunkMetadata;
+ this.valueFilter = filter;
+ }
+
+ @Override
+ public BatchData getAllSatisfiedPageData() throws IOException {
+ return IPageReader.super.getAllSatisfiedPageData();
+ }
+
+ @Override
+ public BatchData getAllSatisfiedPageData(boolean ascending) throws
IOException {
+ TSDataType dataType = chunkMetadata.getDataType();
+ BatchData batchData = BatchDataFactory.createBatchData(dataType,
ascending, false);
+ while (timeValuePairIterator.hasNextTimeValuePair()) {
+ TimeValuePair timeValuePair = timeValuePairIterator.nextTimeValuePair();
+ if (valueFilter == null
+ || valueFilter.satisfy(
+ timeValuePair.getTimestamp(),
timeValuePair.getValue().getValue())) {
+ batchData.putAnObject(timeValuePair.getTimestamp(),
timeValuePair.getValue().getValue());
+ }
+ }
+ return batchData.flip();
+ }
+
+ @Override
+ public Statistics getStatistics() {
+ return chunkMetadata.getStatistics();
+ }
+
+ @Override
+ public Statistics getStatistics(int index) {
+ return chunkMetadata.getStatistics(index);
+ }
+
+ @Override
+ public void setFilter(Filter filter) {
+ if (valueFilter == null) {
+ this.valueFilter = filter;
+ } else {
+ valueFilter = new AndFilter(this.valueFilter, filter);
+ }
+ }
+
+ @Override
+ public boolean isModified() {
+ return false;
+ }
+}
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 12c0436..416f47b 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
@@ -42,8 +42,8 @@ import org.apache.iotdb.tsfile.read.common.BatchData;
import org.apache.iotdb.tsfile.read.common.BatchDataFactory;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
import org.apache.iotdb.tsfile.read.filter.basic.UnaryFilter;
+import org.apache.iotdb.tsfile.read.reader.IAlignedPageReader;
import org.apache.iotdb.tsfile.read.reader.IPageReader;
-import org.apache.iotdb.tsfile.read.reader.page.AlignedPageReader;
import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
import java.io.IOException;
@@ -653,8 +653,8 @@ public class SeriesReader {
if (firstPageReader == null) {
return null;
}
- if (!(firstPageReader.isVectorPageReader())) {
- throw new IOException("Can only get statistics by index from
VectorPageReader");
+ if (!(firstPageReader.isAlignedPageReader())) {
+ throw new IOException("Can only get statistics by index from
AlignedPageReader");
}
return firstPageReader.getStatistics(index);
}
@@ -1087,8 +1087,8 @@ public class SeriesReader {
this.isSeq = isSeq;
}
- public boolean isVectorPageReader() {
- return data instanceof AlignedPageReader;
+ public boolean isAlignedPageReader() {
+ return data instanceof IAlignedPageReader;
}
Statistics getStatistics() {
@@ -1096,10 +1096,10 @@ public class SeriesReader {
}
Statistics getStatistics(int index) throws IOException {
- if (!(data instanceof AlignedPageReader)) {
- throw new IOException("Can only get statistics by index from
VectorPageReader");
+ if (!(data instanceof IAlignedPageReader)) {
+ throw new IOException("Can only get statistics by index from
AlignedPageReader");
}
- return ((AlignedPageReader) data).getStatistics(index);
+ return ((IAlignedPageReader) data).getStatistics(index);
}
BatchData getAllSatisfiedPageData(boolean ascending) throws IOException {
diff --git
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IAlignedPageReader.java
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IAlignedPageReader.java
new file mode 100644
index 0000000..c06a2a7
--- /dev/null
+++
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/IAlignedPageReader.java
@@ -0,0 +1,26 @@
+/*
+ * 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.tsfile.read.reader;
+
+import org.apache.iotdb.tsfile.file.metadata.statistics.Statistics;
+
+public interface IAlignedPageReader {
+
+ Statistics getStatistics(int index);
+}
diff --git
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
index 0a92642..136092e 100644
---
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
+++
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/reader/page/AlignedPageReader.java
@@ -27,6 +27,7 @@ import org.apache.iotdb.tsfile.read.common.BatchDataFactory;
import org.apache.iotdb.tsfile.read.common.TimeRange;
import org.apache.iotdb.tsfile.read.filter.basic.Filter;
import org.apache.iotdb.tsfile.read.filter.operator.AndFilter;
+import org.apache.iotdb.tsfile.read.reader.IAlignedPageReader;
import org.apache.iotdb.tsfile.read.reader.IPageReader;
import org.apache.iotdb.tsfile.utils.TsPrimitiveType;
@@ -35,7 +36,7 @@ import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.List;
-public class AlignedPageReader implements IPageReader {
+public class AlignedPageReader implements IPageReader, IAlignedPageReader {
private final TimePageReader timePageReader;
private final List<ValuePageReader> valuePageReaderList;
@@ -124,6 +125,7 @@ public class AlignedPageReader implements IPageReader {
: timePageReader.getStatistics();
}
+ @Override
public Statistics getStatistics(int index) {
ValuePageReader valuePageReader = valuePageReaderList.get(index);
return valuePageReader == null ? null : valuePageReader.getStatistics();