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();

Reply via email to