This is an automated email from the ASF dual-hosted git repository. jackietien pushed a commit to branch UnseqImprove in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 7ce3c9cd813894a36ed6acc86067080b5b09036c Author: JackieTien97 <[email protected]> AuthorDate: Thu Oct 8 21:09:58 2020 +0800 init --- .../iotdb/db/query/reader/series/SeriesReader.java | 178 +++++++++++++++------ .../iotdb/tsfile/file/metadata/ChunkMetadata.java | 11 ++ .../tsfile/file/metadata/TimeseriesMetadata.java | 11 ++ 3 files changed, 155 insertions(+), 45 deletions(-) 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 e74d4b9..2054f87 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 @@ -95,7 +95,8 @@ public class SeriesReader { * page cache */ private VersionPageReader firstPageReader; - private PriorityQueue<VersionPageReader> cachedPageReaders; + private final List<VersionPageReader> seqPageReaders = new LinkedList<>(); + private PriorityQueue<VersionPageReader> unseqPageReaders; /* * point cache @@ -133,7 +134,7 @@ public class SeriesReader { timeSeriesMetadata -> orderUtils.getOrderTime(timeSeriesMetadata.getStatistics()))); cachedChunkMetadata = new PriorityQueue<>(orderUtils.comparingLong( chunkMetadata -> orderUtils.getOrderTime(chunkMetadata.getStatistics()))); - cachedPageReaders = new PriorityQueue<>(orderUtils.comparingLong( + unseqPageReaders = new PriorityQueue<>(orderUtils.comparingLong( versionPageReader -> orderUtils.getOrderTime(versionPageReader.getStatistics()))); } @@ -163,7 +164,7 @@ public class SeriesReader { timeSeriesMetadata -> orderUtils.getOrderTime(timeSeriesMetadata.getStatistics()))); cachedChunkMetadata = new PriorityQueue<>(orderUtils.comparingLong( chunkMetadata -> orderUtils.getOrderTime(chunkMetadata.getStatistics()))); - cachedPageReaders = new PriorityQueue<>(orderUtils.comparingLong( + unseqPageReaders = new PriorityQueue<>(orderUtils.comparingLong( versionPageReader -> orderUtils.getOrderTime(versionPageReader.getStatistics()))); } @@ -173,12 +174,12 @@ public class SeriesReader { boolean hasNextFile() throws IOException { - if (!cachedPageReaders.isEmpty() + if (!unseqPageReaders.isEmpty() || firstPageReader != null || mergeReader.hasNextTimeValuePair()) { throw new IOException( "all cached pages should be consumed first cachedPageReaders.isEmpty() is " - + cachedPageReaders.isEmpty() + + unseqPageReaders.isEmpty() + " firstPageReader != null is " + (firstPageReader != null) + " mergeReader.hasNextTimeValuePair() = " @@ -231,12 +232,12 @@ public class SeriesReader { * overlapped chunks are consumed */ boolean hasNextChunk() throws IOException { - if (!cachedPageReaders.isEmpty() + if (!unseqPageReaders.isEmpty() || firstPageReader != null || mergeReader.hasNextTimeValuePair()) { throw new IOException( "all cached pages should be consumed first cachedPageReaders.isEmpty() is " - + cachedPageReaders.isEmpty() + + unseqPageReaders.isEmpty() + " firstPageReader != null is " + (firstPageReader != null) + " mergeReader.hasNextTimeValuePair() = " @@ -300,6 +301,7 @@ public class SeriesReader { throws IOException { List<ChunkMetadata> chunkMetadataList = FileLoaderUtils .loadChunkMetadataList(timeSeriesMetadata); + chunkMetadataList.forEach(chunkMetadata -> chunkMetadata.setSeq(timeSeriesMetadata.isSeq())); // try to calculate the total number of chunk and time-value points in chunk if (IoTDBDescriptor.getInstance().getConfig().isEnablePerformanceTracing()) { QueryResourceManager queryResourceManager = QueryResourceManager.getInstance(); @@ -381,19 +383,28 @@ public class SeriesReader { /* * first chunk metadata is already unpacked, consume cached pages */ - if (!cachedPageReaders.isEmpty()) { - firstPageReader = cachedPageReaders.poll(); - long endpointTime = orderUtils.getOverlapCheckTime(firstPageReader.getStatistics()); - unpackAllOverlappedTsFilesToTimeSeriesMetadata(endpointTime); - unpackAllOverlappedTimeSeriesMetadataToCachedChunkMetadata(endpointTime, false); - unpackAllOverlappedChunkMetadataToCachedPageReaders(endpointTime, false); + if (!seqPageReaders.isEmpty() && !unseqPageReaders.isEmpty()) { + if (seqPageReaders.get(0).getStatistics().getStartTime() < unseqPageReaders.peek() + .getStatistics().getStartTime()) { + firstPageReader = seqPageReaders.remove(0); + } else { + firstPageReader = unseqPageReaders.poll(); + } + } else if (!seqPageReaders.isEmpty()) { + firstPageReader = seqPageReaders.remove(0); + } else { + firstPageReader = unseqPageReaders.poll(); } + long endpointTime = orderUtils.getOverlapCheckTime(firstPageReader.getStatistics()); + unpackAllOverlappedTsFilesToTimeSeriesMetadata(endpointTime); + unpackAllOverlappedTimeSeriesMetadataToCachedChunkMetadata(endpointTime, false); + unpackAllOverlappedChunkMetadataToCachedPageReaders(endpointTime, false); } - if (firstPageReader != null - && !cachedPageReaders.isEmpty() - && orderUtils - .isOverlapped(firstPageReader.getStatistics(), cachedPageReaders.peek().getStatistics())) { + if (firstPageReader != null && ((!unseqPageReaders.isEmpty() && orderUtils + .isOverlapped(firstPageReader.getStatistics(), unseqPageReaders.peek().getStatistics())) + || (!seqPageReaders.isEmpty() && orderUtils + .isOverlapped(firstPageReader.getStatistics(), seqPageReaders.get(0).getStatistics())))) { /* * next page is overlapped, read overlapped data and cache it */ @@ -407,11 +418,26 @@ public class SeriesReader { } // make sure firstPageReader won't be null while cachedPageReaders has more cached page readers - while (firstPageReader == null && !cachedPageReaders.isEmpty()) { - firstPageReader = cachedPageReaders.poll(); - if (!cachedPageReaders.isEmpty() - && orderUtils.isOverlapped(firstPageReader.getStatistics(), - cachedPageReaders.peek().getStatistics())) { + while (firstPageReader == null && (!seqPageReaders.isEmpty() || !unseqPageReaders.isEmpty())) { + + if (!seqPageReaders.isEmpty() && !unseqPageReaders.isEmpty()) { + if (seqPageReaders.get(0).getStatistics().getStartTime() < unseqPageReaders.peek() + .getStatistics().getStartTime()) { + firstPageReader = seqPageReaders.remove(0); + } else { + firstPageReader = unseqPageReaders.poll(); + } + } else if (!seqPageReaders.isEmpty()) { + firstPageReader = seqPageReaders.remove(0); + } else { + firstPageReader = unseqPageReaders.poll(); + } + + if ((!seqPageReaders.isEmpty() && orderUtils + .isOverlapped(firstPageReader.getStatistics(), seqPageReaders.get(0).getStatistics())) + || (!unseqPageReaders.isEmpty() && orderUtils + .isOverlapped(firstPageReader.getStatistics(), + unseqPageReaders.peek().getStatistics()))) { /* * next page is overlapped, read overlapped data and cache it */ @@ -438,17 +464,22 @@ public class SeriesReader { unpackOneChunkMetaData(firstChunkMetadata); firstChunkMetadata = null; } - if (init && firstPageReader == null && !cachedPageReaders.isEmpty()) { - firstPageReader = cachedPageReaders.poll(); + if (init && firstPageReader == null && !unseqPageReaders.isEmpty()) { + firstPageReader = unseqPageReaders.poll(); } } - private void unpackOneChunkMetaData(ChunkMetadata chunkMetaData) throws IOException { - FileLoaderUtils.loadPageReaderList(chunkMetaData, timeFilter) - .forEach( - pageReader -> - cachedPageReaders.add( - new VersionPageReader(chunkMetaData.getVersion(), pageReader))); + private void unpackOneChunkMetaData(ChunkMetadata chunkMetaData) + throws IOException { + FileLoaderUtils.loadPageReaderList(chunkMetaData, timeFilter).forEach( + pageReader -> { + if (chunkMetaData.isSeq()) { + seqPageReaders.add(new VersionPageReader(chunkMetaData.getVersion(), pageReader, true)); + } else { + unseqPageReaders + .add(new VersionPageReader(chunkMetaData.getVersion(), pageReader, false)); + } + }); } /** @@ -469,14 +500,16 @@ public class SeriesReader { /* * has a non-overlapped page in firstPageReader */ - if (mergeReader.hasNextTimeValuePair()) { + if (mergeReader.hasNextTimeValuePair() + && mergeReader.currentTimeValuePair().getTimestamp() <= firstPageReader.getStatistics() + .getEndTime()) { throw new IOException("overlapped data should be consumed first"); } Statistics firstPageStatistics = firstPageReader.getStatistics(); - return !cachedPageReaders.isEmpty() - && orderUtils.isOverlapped(firstPageStatistics, cachedPageReaders.peek().getStatistics()); + return !unseqPageReaders.isEmpty() + && orderUtils.isOverlapped(firstPageStatistics, unseqPageReaders.peek().getStatistics()); } Statistics currentPageStatistics() { @@ -543,7 +576,14 @@ public class SeriesReader { cachedBatchData = BatchDataFactory.createBatchData(dataType); long currentPageEndPointTime = mergeReader.getCurrentReadStopTime(); - + if (firstPageReader != null) { + currentPageEndPointTime = Math + .min(currentPageEndPointTime, firstPageReader.getStatistics().getEndTime()); + } + if (!seqPageReaders.isEmpty()) { + currentPageEndPointTime = Math + .min(currentPageEndPointTime, seqPageReaders.get(0).getStatistics().getEndTime()); + } while (mergeReader.hasNextTimeValuePair()) { /* @@ -559,7 +599,35 @@ public class SeriesReader { unpackAllOverlappedTimeSeriesMetadataToCachedChunkMetadata( timeValuePair.getTimestamp(), false); unpackAllOverlappedChunkMetadataToCachedPageReaders(timeValuePair.getTimestamp(), false); - unpackAllOverlappedCachedPageReadersToMergeReader(timeValuePair.getTimestamp()); + unpackAllOverlappedUnseqPageReadersToMergeReader(timeValuePair.getTimestamp()); + + /* + * get the latest first point in mergeReader + */ + timeValuePair = mergeReader.currentTimeValuePair(); + + if (firstPageReader != null) { + if (firstPageReader.getStatistics().getEndTime() < timeValuePair.getTimestamp()) { + return cachedBatchData.hasCurrent(); + } else { + mergeReader + .addReader(firstPageReader.getAllSatisfiedPageData(orderUtils.getAscending()) + .getBatchDataIterator(), firstPageReader.version, + orderUtils.getOverlapCheckTime(firstPageReader.getStatistics())); + firstPageReader = null; + } + } + + if (!seqPageReaders.isEmpty()) { + if (seqPageReaders.get(0).getStatistics().getEndTime() < timeValuePair.getTimestamp()) { + return cachedBatchData.hasCurrent(); + } else { + VersionPageReader pageReader = seqPageReaders.remove(0); + mergeReader.addReader(pageReader.getAllSatisfiedPageData(orderUtils.getAscending()) + .getBatchDataIterator(), pageReader.version, + orderUtils.getOverlapCheckTime(pageReader.getStatistics())); + } + } /* * get the latest first point in mergeReader @@ -591,7 +659,7 @@ public class SeriesReader { /* * no cached page readers */ - if (firstPageReader == null && cachedPageReaders.isEmpty()) { + if (firstPageReader == null && unseqPageReaders.isEmpty() && seqPageReaders.isEmpty()) { return; } @@ -599,7 +667,18 @@ public class SeriesReader { * init firstPageReader */ if (firstPageReader == null) { - firstPageReader = cachedPageReaders.poll(); + if (!seqPageReaders.isEmpty() && !unseqPageReaders.isEmpty()) { + if (seqPageReaders.get(0).getStatistics().getStartTime() < unseqPageReaders.peek() + .getStatistics().getStartTime()) { + firstPageReader = seqPageReaders.remove(0); + } else { + firstPageReader = unseqPageReaders.poll(); + } + } else if (!seqPageReaders.isEmpty()) { + firstPageReader = seqPageReaders.remove(0); + } else { + firstPageReader = unseqPageReaders.poll(); + } } long currentPageEndpointTime; @@ -610,18 +689,18 @@ public class SeriesReader { } /* - * put all currently directly overlapped page reader to merge reader + * put all currently directly overlapped unseq page reader to merge reader */ - unpackAllOverlappedCachedPageReadersToMergeReader(currentPageEndpointTime); + unpackAllOverlappedUnseqPageReadersToMergeReader(currentPageEndpointTime); } - private void unpackAllOverlappedCachedPageReadersToMergeReader(long endpointTime) + private void unpackAllOverlappedUnseqPageReadersToMergeReader(long endpointTime) throws IOException { - while (!cachedPageReaders.isEmpty() - && orderUtils.isOverlapped(endpointTime, cachedPageReaders.peek().data.getStatistics())) { - putPageReaderToMergeReader(cachedPageReaders.poll()); + while (!unseqPageReaders.isEmpty() + && orderUtils.isOverlapped(endpointTime, unseqPageReaders.peek().data.getStatistics())) { + putPageReaderToMergeReader(unseqPageReaders.poll()); } - if (firstPageReader != null && + if (firstPageReader != null && !firstPageReader.isSeq && orderUtils.isOverlapped(endpointTime, firstPageReader.getStatistics())) { putPageReaderToMergeReader(firstPageReader); firstPageReader = null; @@ -668,6 +747,7 @@ public class SeriesReader { orderUtils.getNextSeqFileResource(seqFileResource, true), seriesPath, context, getAnyFilter(), allSensors); if (timeseriesMetadata != null) { + timeseriesMetadata.setSeq(true); seqTimeSeriesMetadata.add(timeseriesMetadata); } } @@ -681,6 +761,7 @@ public class SeriesReader { unseqFileResource.remove(0), seriesPath, context, getAnyFilter(), allSensors); if (timeseriesMetadata != null) { timeseriesMetadata.setModified(true); + timeseriesMetadata.setSeq(false); unSeqTimeSeriesMetadata.add(timeseriesMetadata); } } @@ -775,9 +856,12 @@ public class SeriesReader { protected long version; protected IPageReader data; - VersionPageReader(long version, IPageReader data) { + protected boolean isSeq; + + VersionPageReader(long version, IPageReader data, boolean isSeq) { this.version = version; this.data = data; + this.isSeq = isSeq; } Statistics getStatistics() { @@ -795,6 +879,10 @@ public class SeriesReader { boolean isModified() { return data.isModified(); } + + public boolean isSeq() { + return isSeq; + } } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java index da68c93..499bf57 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/ChunkMetadata.java @@ -73,6 +73,9 @@ public class ChunkMetadata implements Accountable { private static final int CHUNK_METADATA_FIXED_RAM_SIZE = 80; + // used for SeriesReader to indicate whether it is a seq/unseq timeseries metadata + private boolean isSeq = true; + private ChunkMetadata() { } @@ -269,4 +272,12 @@ public class ChunkMetadata implements Accountable { this.statistics.mergeStatistics(chunkMetadata.getStatistics()); this.ramSize = calculateRamSize(); } + + public void setSeq(boolean seq) { + isSeq = seq; + } + + public boolean isSeq() { + return isSeq; + } } \ No newline at end of file diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java index 754d48b..ad8fc3f 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/file/metadata/TimeseriesMetadata.java @@ -46,6 +46,9 @@ public class TimeseriesMetadata implements Accountable { private long ramSize; + // used for SeriesReader to indicate whether it is a seq/unseq timeseries metadata + private boolean isSeq = true; + public TimeseriesMetadata() { } @@ -158,4 +161,12 @@ public class TimeseriesMetadata implements Accountable { public long getRamSize() { return ramSize; } + + public void setSeq(boolean seq) { + isSeq = seq; + } + + public boolean isSeq() { + return isSeq; + } }
