From fa993337aa782fd7b230a21df331894e081cc3fb Mon Sep 17 00:00:00 2001 From: "huangjianfeng.118" Date: Thu, 23 Jul 2026 19:27:17 +0800 Subject: [PATCH 1/3] Add Parquet repdef preload window option plumbing --- bolt/dwio/common/Options.h | 15 +++++++++++++-- bolt/dwio/parquet/reader/PageReader.h | 5 +++++ bolt/dwio/parquet/reader/ParquetData.cpp | 3 +++ bolt/dwio/parquet/reader/ParquetData.h | 6 ++++++ bolt/dwio/parquet/reader/ParquetReader.cpp | 1 + 5 files changed, 28 insertions(+), 2 deletions(-) diff --git a/bolt/dwio/common/Options.h b/bolt/dwio/common/Options.h index b3d4447ac..0412195a3 100644 --- a/bolt/dwio/common/Options.h +++ b/bolt/dwio/common/Options.h @@ -175,6 +175,7 @@ class RowReaderOptions { bool enableDictionaryFilter_ = false; int32_t decodeRepDefPageCount_{10}; + int32_t parquetRepDefPreloadWindowCount_{0}; int32_t parquetRepDefMemoryLimit_{16UL << 20}; bool useColumnNamesForColumnMapping_{false}; @@ -483,6 +484,14 @@ class RowReaderOptions { return decodeRepDefPageCount_; } + void setParquetRepDefPreloadWindowCount(int32_t windowCount) { + parquetRepDefPreloadWindowCount_ = windowCount < 0 ? 0 : windowCount; + } + + int32_t getParquetRepDefPreloadWindowCount() const { + return parquetRepDefPreloadWindowCount_; + } + void setParquetRepDefMemoryLimit(int32_t memlimit) { parquetRepDefMemoryLimit_ = memlimit; } @@ -556,8 +565,10 @@ class RowReaderOptions { << ", "; ss << "fileName_=" << fileName_ << ", "; ss << "fileId_=" << fileId_ << ", "; - ss << "enableDictionaryFilter_=" << enableDictionaryFilter_; - ss << "decodeRepDefPageCount_=" << decodeRepDefPageCount_; + ss << "enableDictionaryFilter_=" << enableDictionaryFilter_ << ", "; + ss << "decodeRepDefPageCount_=" << decodeRepDefPageCount_ << ", "; + ss << "parquetRepDefPreloadWindowCount_=" + << parquetRepDefPreloadWindowCount_ << ", "; ss << "parquetRepDefMemoryLimit_=" << parquetRepDefMemoryLimit_ << ", "; ss << "maxBatchBytes_=" << maxBatchBytes_ << ", "; ss << "useColumnNamesForColumnMapping_=" diff --git a/bolt/dwio/parquet/reader/PageReader.h b/bolt/dwio/parquet/reader/PageReader.h index 51038529e..713141a83 100644 --- a/bolt/dwio/parquet/reader/PageReader.h +++ b/bolt/dwio/parquet/reader/PageReader.h @@ -247,6 +247,10 @@ class PageReader { decodeRepDefPageCount_ = count; } + void setParquetRepDefPreloadWindowCount(int32_t count) { + repDefPreloadWindowCount_ = count < 0 ? 0 : count; + } + void setParquetRepDefMemoryLimit(int32_t memlimit) { repDefMemoryLimit_ = memlimit; } @@ -629,6 +633,7 @@ class PageReader { int32_t repDefMemoryLimit_{16L << 20}; int64_t totalRefDefBytes_{0}; int32_t decodeRepDefPageCount_{10}; + int32_t repDefPreloadWindowCount_{0}; dwio::common::RuntimeStatistics* statis_{nullptr}; // Tracks output count for the current physical page. -1 means there is no diff --git a/bolt/dwio/parquet/reader/ParquetData.cpp b/bolt/dwio/parquet/reader/ParquetData.cpp index 750eb0693..bf8e16564 100644 --- a/bolt/dwio/parquet/reader/ParquetData.cpp +++ b/bolt/dwio/parquet/reader/ParquetData.cpp @@ -49,6 +49,7 @@ std::unique_ptr ParquetParams::toFormatData( schemaHelper_, enableDictionaryFilter_, decodeRepDefPageCount_, + parquetRepDefPreloadWindowCount_, parquetRepDefMemoryLimit_); } @@ -373,6 +374,8 @@ dwio::common::PositionProvider ParquetData::seekToRowGroup(int64_t index) { metadata.total_compressed_size, statis_); reader_->setDecodeRepDefPageCount(decodeRepDefPageCount_); + reader_->setParquetRepDefPreloadWindowCount( + parquetRepDefPreloadWindowCount_); reader_->setParquetRepDefMemoryLimit(parquetRepDefMemoryLimit_); if (columnChunkMeta.__isset.crypto_metadata) { diff --git a/bolt/dwio/parquet/reader/ParquetData.h b/bolt/dwio/parquet/reader/ParquetData.h index 0bd7ccecc..31bfe5957 100644 --- a/bolt/dwio/parquet/reader/ParquetData.h +++ b/bolt/dwio/parquet/reader/ParquetData.h @@ -61,6 +61,7 @@ class ParquetParams : public dwio::common::FormatParams { const SchemaHelper& schemaHelper, bool enableDictionaryFilter, int32_t decodeRepDefPageCount, + int32_t parquetRepDefPreloadWindowCount, int32_t parquetRepDefMemoryLimit) : FormatParams(pool, stats), metaData_(metaData), @@ -69,6 +70,7 @@ class ParquetParams : public dwio::common::FormatParams { schemaHelper_(schemaHelper), enableDictionaryFilter_(enableDictionaryFilter), decodeRepDefPageCount_(decodeRepDefPageCount), + parquetRepDefPreloadWindowCount_(parquetRepDefPreloadWindowCount), parquetRepDefMemoryLimit_(parquetRepDefMemoryLimit) {} std::unique_ptr toFormatData( @@ -90,6 +92,7 @@ class ParquetParams : public dwio::common::FormatParams { const SchemaHelper& schemaHelper_; const bool enableDictionaryFilter_; const int32_t decodeRepDefPageCount_; + const int32_t parquetRepDefPreloadWindowCount_; const int32_t parquetRepDefMemoryLimit_; }; @@ -105,6 +108,7 @@ class ParquetData : public dwio::common::FormatData { const SchemaHelper& schemaHelper, bool enableDictionaryFilter, int32_t decodeRepDefPageCount, + int32_t parquetRepDefPreloadWindowCount, int32_t parquetRepDefMemoryLimit) : pool_(pool), type_(std::static_pointer_cast(type)), @@ -117,6 +121,7 @@ class ParquetData : public dwio::common::FormatData { schemaHelper_(schemaHelper), enableDictionaryFilter_(enableDictionaryFilter), decodeRepDefPageCount_(decodeRepDefPageCount), + parquetRepDefPreloadWindowCount_(parquetRepDefPreloadWindowCount), parquetRepDefMemoryLimit_(parquetRepDefMemoryLimit) { rowGroupOffsets_ = std::vector(rowGroups_.size()); rowGroupOffsets_[0] = 0; @@ -331,6 +336,7 @@ class ParquetData : public dwio::common::FormatData { const SchemaHelper& schemaHelper_; const bool enableDictionaryFilter_; const int32_t decodeRepDefPageCount_; + const int32_t parquetRepDefPreloadWindowCount_; const int32_t parquetRepDefMemoryLimit_; }; diff --git a/bolt/dwio/parquet/reader/ParquetReader.cpp b/bolt/dwio/parquet/reader/ParquetReader.cpp index 05fe426e5..dd7249c54 100644 --- a/bolt/dwio/parquet/reader/ParquetReader.cpp +++ b/bolt/dwio/parquet/reader/ParquetReader.cpp @@ -1536,6 +1536,7 @@ class ParquetRowReader::Impl { schemaHelper_, options_.isDictionaryFilterEnabled(), options_.getDecodeRepDefPageCount(), + options_.getParquetRepDefPreloadWindowCount(), options_.getParquetRepDefMemoryLimit()); if (auto selector = options_.getSelector()) { From 1837b58e39edbdd21d9f02eb63a7809d59fd69d1 Mon Sep 17 00:00:00 2001 From: "huangjianfeng.118" Date: Thu, 23 Jul 2026 19:34:17 +0800 Subject: [PATCH 2/3] Add Parquet repdef preload window mode --- bolt/dwio/parquet/reader/PageReader.cpp | 421 ++++++++++++++++++++++-- bolt/dwio/parquet/reader/PageReader.h | 15 +- 2 files changed, 412 insertions(+), 24 deletions(-) diff --git a/bolt/dwio/parquet/reader/PageReader.cpp b/bolt/dwio/parquet/reader/PageReader.cpp index bbec00ebb..a061c8c3b 100644 --- a/bolt/dwio/parquet/reader/PageReader.cpp +++ b/bolt/dwio/parquet/reader/PageReader.cpp @@ -33,6 +33,13 @@ #include #include +#include +#include +#include +#include + +#include "bolt/common/base/RuntimeMetrics.h" +#include "bolt/common/time/Timer.h" #include "bolt/dwio/common/BufferUtil.h" #include "bolt/dwio/common/ColumnVisitors.h" #include "bolt/dwio/parquet/reader/Decompression.h" @@ -49,6 +56,34 @@ namespace bytedance::bolt::parquet { using thrift::Encoding; using thrift::PageHeader; +namespace { + +uint64_t configuredRepDefPreloadWindowCount() { + const char* configured = + std::getenv("BOLT_PARQUET_REPDEF_PRELOAD_WINDOW_COUNT"); + if (configured == nullptr || *configured == '\0') { + return 0; + } + + char* end = nullptr; + const auto parsed = std::strtoull(configured, &end, 10); + if (end == configured) { + return 0; + } + return parsed; +} + +void addPageReaderStat(const std::string& name, int64_t value) { + addThreadLocalRuntimeStat(name, RuntimeCounter(value)); +} + +void addPageReaderNanos(const std::string& name, uint64_t value) { + addThreadLocalRuntimeStat( + name, RuntimeCounter(value, RuntimeCounter::Unit::kNanos)); +} + +} // namespace + void PageReader::seekToPage(int64_t row, bool keepRepDefRawData) { defineDecoder_.reset(); repeatDecoder_.reset(); @@ -68,21 +103,10 @@ void PageReader::seekToPage(int64_t row, bool keepRepDefRawData) { numRowsInPage_ = 0; break; } - // For nested (non top level) columns, 'preloadRepDefs' only fully - // decodes the first 'decodeRepDefPageCount_' pages and keeps the - // remaining pages as raw rep/def bytes inside 'preloadedRepDefs_'. - // When a leaf-level filter pushdown triggers a skip()/seekToPage() - // that walks past the sampled boundary, the consumer side - // ('setPageRowInfo' inside 'prepareDataPageV1') indexes - // 'numLeavesInPage_[pageIndex_]' and would go out of bounds. - // - // Materialise just enough preloaded rep/def bytes BEFORE entering - // the next page so that 'setPageRowInfo' remains a pure lookup. - // The loop only fires when we are about to cross the sampled - // boundary and there is something left to decode, mirroring (and - // complementing) the "ahead by one page" invariant established at - // the exit of 'decodeRepDefs'. - if (hasChunkRepDefs_ && !isTopLevel_ && maxRepeat_ > 0) { + if (row != kRepDefOnly && hasChunkRepDefs_ && !isTopLevel_ && + maxRepeat_ > 0) { + // In non-window mode we still need the existing sampled-preload guard + // to materialize leaf counts before crossing the sampling boundary. while (pageIndex_ + 1 >= static_cast(numLeavesInPage_.size()) && !preloadedRepDefs_.empty()) { loadMoreRepDefs(); @@ -771,6 +795,7 @@ int32_t parquetTypeBytes(thrift::Type::type type) { } // namespace void PageReader::preloadPageRepDefs(const bool keepRepDefRawData) { + const auto startNs = getCurrentTimeNano(); seekToPage(kRepDefOnly, keepRepDefRawData); if (!keepRepDefRawData) { BOLT_CHECK_GE( @@ -807,11 +832,21 @@ void PageReader::preloadPageRepDefs(const bool keepRepDefRawData) { leafNullsSize_ += numLeaves; numLeavesInPage_.push_back(numLeaves); } + addPageReaderStat("pageReaderPreloadPageRepDefsCount", 1); + addPageReaderStat("pageReaderPreloadPageRepDefsLevels", numRepDefsInPage_); + if (keepRepDefRawData) { + addPageReaderStat("pageReaderPreloadPageRepDefsKeepRawCount", 1); + } + addPageReaderNanos( + "pageReaderPreloadPageRepDefsTimeNs", getCurrentTimeNano() - startNs); return; } void PageReader::preloadRepDefs() { + const auto startNs = getCurrentTimeNano(); hasChunkRepDefs_ = true; + windowRepDefPageData_.clear(); + windowRepDefsRemaining_ = 0; bool startWithDictinoaryPage = cryptoCtx_.startDecryptWithDictionaryPage; if (maxRepeat_ > 0 && maxDefine_ > 0) { const int32_t samplePages = decodeRepDefPageCount_; @@ -830,7 +865,7 @@ void PageReader::preloadRepDefs() { 1; auto avgRefDefBytes = totalRefDefBytes_ / samplePages + 1; auto avgPageLen = pageStart_ / samplePages + 1; - LOG(INFO) << __FUNCTION__ << "[" << this + VLOG(1) << __FUNCTION__ << "[" << this << "] : def/rep levels = " << definitionLevels_.size() << ", avgPageLen = " << avgPageLen << ", chunkSize = " << chunkSize_ @@ -871,18 +906,327 @@ void PageReader::preloadRepDefs() { totalRefDefBytes_ = 0; pageOrdinal_ = 0; cryptoCtx_.startDecryptWithDictionaryPage = startWithDictinoaryPage; + addPageReaderStat("pageReaderPreloadRepDefsCount", 1); + addPageReaderNanos( + "pageReaderPreloadRepDefsTimeNs", getCurrentTimeNano() - startNs); +} + +void PageReader::preloadRepDefsRawOnly() { + const auto startNs = getCurrentTimeNano(); + hasChunkRepDefs_ = true; + const bool startWithDictionaryPage = + cryptoCtx_.startDecryptWithDictionaryPage; + totalRefDefBytes_ = 0; + numLeavesInPage_.clear(); + windowRepeatDecoder_.reset(); + windowDefineDecoder_.reset(); + windowRepDefsRemaining_ = 0; + windowRepDefPageData_.clear(); + + while (pageStart_ < chunkSize_) { + preloadedRepDefs_.emplace_back(raw_vector(&pool_)); + preloadPageRepDefs(true); + if (preloadedRepDefs_.back().empty()) { + preloadedRepDefs_.pop_back(); + } else { + numLeavesInPage_.push_back( + countLeavesInRepDefPage(preloadedRepDefs_.back())); + } + } + + std::vector rewind = {0}; + pageStart_ = 0; + dwio::common::PositionProvider position(rewind); + inputStream_->seekToPosition(position); + bufferStart_ = bufferEnd_ = nullptr; + rowOfPage_ = 0; + numRowsInPage_ = 0; + pageData_ = nullptr; + totalRefDefBytes_ = 0; + pageOrdinal_ = 0; + cryptoCtx_.startDecryptWithDictionaryPage = startWithDictionaryPage; + addPageReaderStat("pageReaderPreloadRepDefsRawOnlyCount", 1); + addPageReaderNanos( + "pageReaderPreloadRepDefsRawOnlyTimeNs", + getCurrentTimeNano() - startNs); +} + +int32_t PageReader::countLeavesInRepDefPage( + const raw_vector& repDefData) { + if (repDefData.empty()) { + return 0; + } + + const char* rawData = repDefData.data(); + const char* pageEnd = repDefData.data() + repDefData.size(); + const auto numRepDefsInPage = readField(rawData); + if (numRepDefsInPage <= 0) { + return 0; + } + + const auto repeatLength = readField(rawData); + BOLT_CHECK_GE(repeatLength, 0); + BOLT_CHECK_LE(rawData + repeatLength, pageEnd); + rawData += repeatLength; + + const auto defineLength = readField(rawData); + BOLT_CHECK_LE(rawData + defineLength, pageEnd); + auto defineDecoder = ::arrow::util::RleDecoder( + reinterpret_cast(rawData), + defineLength, + ::arrow::bit_util::NumRequiredBits(maxDefine_)); + + constexpr int32_t kMaxRepDefDecodeBatch = 64 * 1024; + raw_vector definitionScratch(&pool_); + raw_vector nullScratch(&pool_); + definitionScratch.resize(std::min(numRepDefsInPage, kMaxRepDefDecodeBatch)); + + int32_t decoded = 0; + int32_t numLeaves = 0; + while (decoded < numRepDefsInPage) { + const auto batchSize = + std::min(numRepDefsInPage - decoded, kMaxRepDefDecodeBatch); + defineDecoder.GetBatch(definitionScratch.data(), batchSize); + nullScratch.resize(bits::nwords(batchSize)); + arrow::ValidityBitmapInputOutput bits; + bits.values_read_upper_bound = batchSize; + bits.values_read = 0; + bits.null_count = 0; + bits.valid_bits = reinterpret_cast(nullScratch.data()); + bits.valid_bits_offset = 0; + DefLevelsToBitmap(definitionScratch.data(), batchSize, leafInfo_, &bits); + numLeaves += bits.values_read; + decoded += batchSize; + } + return numLeaves; +} + +void PageReader::compactConsumedLeafNulls() { + constexpr int32_t WordBits = 64; + int64_t erasedBits = erasedLeafNullWords_ * WordBits; + BOLT_CHECK_LE(numLeafNullsConsumed_, leafNullsSize_ + erasedBits); + if (numLeafNullsConsumed_ - erasedBits <= WordBits) { + return; + } + + auto consumedWords = bits::nwords(numLeafNullsConsumed_ - erasedBits) - 1; + auto totalNullWords = bits::nwords(leafNullsSize_); + auto unConsumedNullWords = totalNullWords - consumedWords; + if (unConsumedNullWords > 0) { + uint64_t* rawData = leafNulls_.data(); + size_t copyBytes = unConsumedNullWords * sizeof(uint64_t); + +#if (defined(__x86_64__) || defined(__i386__)) && defined(__GLIBC__) && \ + (__GLIBC__ * 100 + __GLIBC_MINOR__ < 233) +#if defined(__GNUC__) && !defined(__clang__) +#pragma GCC diagnostic push +#pragma GCC diagnostic ignored "-Wrestrict" +#elif defined(__clang__) +#pragma clang diagnostic push +#pragma clang diagnostic ignored "-Wunknown-warning-option" +#pragma clang diagnostic ignored "-Wrestrict" +#endif + if (unConsumedNullWords < consumedWords) { + memcpy(rawData, rawData + consumedWords, copyBytes); + } else { + raw_vector tmpBuffer; + tmpBuffer.resize(unConsumedNullWords); + uint64_t* tmpDest = tmpBuffer.data(); + memcpy(tmpDest, rawData + consumedWords, copyBytes); + memcpy(rawData, tmpDest, copyBytes); + } +#if defined(__GNUC__) && !defined(__clang__) +#pragma GCC diagnostic pop +#elif defined(__clang__) +#pragma clang diagnostic pop +#endif +#else + memmove(rawData, rawData + consumedWords, copyBytes); +#endif + } + leafNulls_.resize(unConsumedNullWords); + leafNullsSize_ -= consumedWords * WordBits; + erasedLeafNullWords_ += consumedWords; + BOLT_CHECK_EQ(leafNulls_.size(), bits::nwords(leafNullsSize_)); +} + +void PageReader::compactWindowRepDefs() { + if (repDefEnd_ == 0) { + return; + } + + BOLT_CHECK_LE(repDefEnd_, definitionLevels_.size()); + BOLT_CHECK_EQ(definitionLevels_.size(), repetitionLevels_.size()); + const auto keptLevels = + static_cast(definitionLevels_.size()) - repDefEnd_; + if (keptLevels > 0) { + const auto copySize = keptLevels * sizeof(int16_t); + memmove( + definitionLevels_.data(), + definitionLevels_.data() + repDefEnd_, + copySize); + memmove( + repetitionLevels_.data(), + repetitionLevels_.data() + repDefEnd_, + copySize); + } + definitionLevels_.resize(keptLevels); + repetitionLevels_.resize(keptLevels); + repDefBegin_ = 0; + repDefEnd_ = 0; +} + +bool PageReader::loadNextWindowRepDefPage() { + if (windowRepDefsRemaining_ > 0) { + return true; + } + + windowRepeatDecoder_.reset(); + windowDefineDecoder_.reset(); + while (!preloadedRepDefs_.empty()) { + windowRepDefPageData_ = std::move(preloadedRepDefs_.front()); + preloadedRepDefs_.pop_front(); + if (windowRepDefPageData_.empty()) { + continue; + } + + const char* rawData = windowRepDefPageData_.data(); + const char* pageEnd = + windowRepDefPageData_.data() + windowRepDefPageData_.size(); + windowRepDefsRemaining_ = readField(rawData); + if (windowRepDefsRemaining_ <= 0) { + continue; + } + + const auto repeatLength = readField(rawData); + BOLT_CHECK_GE(repeatLength, 0); + BOLT_CHECK_LE(rawData + repeatLength, pageEnd); + windowRepeatDecoder_ = std::make_unique<::arrow::util::RleDecoder>( + reinterpret_cast(rawData), + repeatLength, + ::arrow::bit_util::NumRequiredBits(maxRepeat_)); + rawData += repeatLength; + + const auto defineLength = readField(rawData); + BOLT_CHECK_LE(rawData + defineLength, pageEnd); + windowDefineDecoder_ = std::make_unique<::arrow::util::RleDecoder>( + reinterpret_cast(rawData), + defineLength, + ::arrow::bit_util::NumRequiredBits(maxDefine_)); + return true; + } + return false; +} + +int32_t PageReader::countTopLevelRepDefs(int32_t begin, int32_t end) const { + if (maxRepeat_ == 0) { + return std::max(end - begin, 0); + } + int32_t topLevelRows = 0; + for (int32_t i = begin; i < end; ++i) { + topLevelRows += repetitionLevels_[i] == 0 ? 1 : 0; + } + return topLevelRows; +} + +void PageReader::loadWindowRepDefs(int32_t targetTopLevelRows) { + const auto startNs = getCurrentTimeNano(); + if (targetTopLevelRows <= 0) { + return; + } + + constexpr int32_t kMaxRepDefDecodeBatch = 64 * 1024; + compactConsumedLeafNulls(); + auto topLevelRows = + countTopLevelRepDefs(repDefBegin_, definitionLevels_.size()); + int32_t decodeBatchCount = 0; + while (topLevelRows < targetTopLevelRows && loadNextWindowRepDefPage()) { + const auto batchSize = + std::min(windowRepDefsRemaining_, kMaxRepDefDecodeBatch); + if (batchSize <= 0) { + continue; + } + + const auto begin = static_cast(definitionLevels_.size()); + const auto end = begin + batchSize; + repetitionLevels_.resize(end); + BOLT_CHECK_NOT_NULL( + windowRepeatDecoder_, "Missing incremental repetition decoder"); + windowRepeatDecoder_->GetBatch(repetitionLevels_.data() + begin, batchSize); + + definitionLevels_.resize(end); + BOLT_CHECK_NOT_NULL( + windowDefineDecoder_, "Missing incremental definition decoder"); + windowDefineDecoder_->GetBatch(definitionLevels_.data() + begin, batchSize); + + windowRepDefsRemaining_ -= batchSize; + leafNulls_.resize(bits::nwords(leafNullsSize_ + batchSize)); + const auto numLeaves = getLengthsAndNulls( + LevelMode::kNulls, + leafInfo_, + begin, + end, + batchSize, + nullptr, + leafNulls_.data(), + leafNullsSize_); + leafNullsSize_ += numLeaves; + topLevelRows += countTopLevelRepDefs(begin, end); + ++decodeBatchCount; + } + addPageReaderStat("pageReaderLoadWindowRepDefsCount", 1); + addPageReaderStat( + "pageReaderLoadWindowRepDefsBatches", decodeBatchCount); + addPageReaderNanos( + "pageReaderLoadWindowRepDefsTimeNs", + getCurrentTimeNano() - startNs); } void PageReader::decodeRepDefs(int32_t numTopLevelRows) { - if (definitionLevels_.empty() && maxDefine_ > 0) { - preloadRepDefs(); + const auto startNs = getCurrentTimeNano(); + uint64_t preloadTimeNs = 0; + uint64_t findTopLevelTimeNs = 0; + uint64_t loadMoreTimeNs = 0; + int64_t findTopLevelLevels = 0; + const auto preloadWindowCount = repDefPreloadWindowCount_ > 0 + ? static_cast(repDefPreloadWindowCount_) + : configuredRepDefPreloadWindowCount(); + const bool windowPreload = + preloadWindowCount > 0 && maxRepeat_ > 0 && maxDefine_ > 0; + if (windowPreload) { + compactWindowRepDefs(); + } + if (!hasChunkRepDefs_ && maxDefine_ > 0) { + const auto preloadStartNs = getCurrentTimeNano(); + if (windowPreload) { + preloadRepDefsRawOnly(); + } else { + preloadRepDefs(); + } + preloadTimeNs += getCurrentTimeNano() - preloadStartNs; } repDefBegin_ = repDefEnd_; + const auto requiredTopLevelRows = numTopLevelRows + 1; + const auto targetTopLevelRows = windowPreload + ? std::max( + requiredTopLevelRows, + static_cast(std::min( + static_cast(std::numeric_limits::max()), + static_cast(numTopLevelRows) * preloadWindowCount + + 1))) + : requiredTopLevelRows; + if (windowPreload) { + const auto loadMoreStartNs = getCurrentTimeNano(); + loadWindowRepDefs(targetTopLevelRows); + loadMoreTimeNs += getCurrentTimeNano() - loadMoreStartNs; + } int32_t numLevels = definitionLevels_.size(); int32_t topFound = 0; int32_t i = repDefBegin_; + int32_t loadMoreCount = 0; auto foundTopLevel = [&]() { + const auto scanBegin = i; for (; i < numLevels; ++i) { if (repetitionLevels_[i] == 0) { ++topFound; @@ -891,32 +1235,63 @@ void PageReader::decodeRepDefs(int32_t numTopLevelRows) { } } } + auto scannedLevels = i - scanBegin; + if (i < numLevels && topFound == numTopLevelRows + 1) { + ++scannedLevels; + } repDefEnd_ = i; + return scannedLevels; }; if (maxRepeat_ > 0) { - foundTopLevel(); + const auto findStartNs = getCurrentTimeNano(); + findTopLevelLevels += foundTopLevel(); + findTopLevelTimeNs += getCurrentTimeNano() - findStartNs; } else { repDefEnd_ = i + numTopLevelRows; } if (maxRepeat_ > 0 && maxDefine_ > 0) { // definitionLevels_ has been consumed, decode more if any while (repDefEnd_ == numLevels && topFound < numTopLevelRows + 1 && - !preloadedRepDefs_.empty()) { - loadMoreRepDefs(); + (!preloadedRepDefs_.empty() || windowRepDefsRemaining_ > 0)) { + const auto loadMoreStartNs = getCurrentTimeNano(); + if (windowPreload) { + loadWindowRepDefs(targetTopLevelRows); + } else { + loadMoreRepDefs(); + } + ++loadMoreCount; numLevels = definitionLevels_.size(); i = repDefEnd_; - foundTopLevel(); + findTopLevelLevels += foundTopLevel(); + loadMoreTimeNs += getCurrentTimeNano() - loadMoreStartNs; BOLT_CHECK(topFound == numTopLevelRows + 1 || repDefEnd_ == numLevels); } // after topFound done, left decoded rep/def is less than 1 page, decode // more - if (!numLeavesInPage_.empty() && + if (!windowPreload && !numLeavesInPage_.empty() && numLevels - repDefEnd_ < numLeavesInPage_.back() && topFound == numTopLevelRows + 1 && !preloadedRepDefs_.empty()) { + const auto loadMoreStartNs = getCurrentTimeNano(); loadMoreRepDefs(); + ++loadMoreCount; + loadMoreTimeNs += getCurrentTimeNano() - loadMoreStartNs; } } + const auto elapsedNs = getCurrentTimeNano() - startNs; + addPageReaderStat("pageReaderDecodeRepDefsRows", numTopLevelRows); + addPageReaderStat("pageReaderDecodeRepDefsLevels", repDefEnd_ - repDefBegin_); + addPageReaderStat("pageReaderDecodeRepDefsLoadMoreCount", loadMoreCount); + addPageReaderStat( + "pageReaderDecodeRepDefsFindTopLevelLevels", findTopLevelLevels); + addPageReaderNanos("pageReaderDecodeRepDefsTotalTimeNs", elapsedNs); + addPageReaderNanos("pageReaderDecodeRepDefsPreloadTimeNs", preloadTimeNs); + addPageReaderNanos( + "pageReaderDecodeRepDefsFindTopLevelTimeNs", findTopLevelTimeNs); + addPageReaderNanos("pageReaderDecodeRepDefsLoadMoreTimeNs", loadMoreTimeNs); + // Keep the legacy aggregate for compatibility while final reporting moves to + // the non-overlapping breakdown above. + addPageReaderNanos("pageReaderDecodeRepDefsTimeNs", elapsedNs); } void PageReader::loadMoreRepDefs() { diff --git a/bolt/dwio/parquet/reader/PageReader.h b/bolt/dwio/parquet/reader/PageReader.h index 713141a83..842b4f953 100644 --- a/bolt/dwio/parquet/reader/PageReader.h +++ b/bolt/dwio/parquet/reader/PageReader.h @@ -93,6 +93,7 @@ class PageReader { definitionLevels_(&pool_), repetitionLevels_(&pool_), nullConcatenation_(pool_), + windowRepDefPageData_(&pool_), statis_(statis) { type_->makeLevelInfo(leafInfo_); pageOrdinal_ = 0; @@ -113,7 +114,8 @@ class PageReader { chunkSize_(chunkSize), definitionLevels_(&pool_), repetitionLevels_(&pool_), - nullConcatenation_(pool_) { + nullConcatenation_(pool_), + windowRepDefPageData_(&pool_) { pageOrdinal_ = 0; } @@ -473,6 +475,13 @@ class PageReader { void preloadPageRepDefs(const bool keepRepDefRawData); void loadMoreRepDefs(); + void preloadRepDefsRawOnly(); + int32_t countLeavesInRepDefPage(const raw_vector& repDefData); + void compactConsumedLeafNulls(); + void compactWindowRepDefs(); + void loadWindowRepDefs(int32_t targetTopLevelRows); + bool loadNextWindowRepDefPage(); + int32_t countTopLevelRepDefs(int32_t begin, int32_t end) const; void decodeRepDefsFromBuffer(); memory::MemoryPool& pool_; @@ -630,6 +639,10 @@ class PageReader { // preload undecoded RepDefs std::list> preloadedRepDefs_; + raw_vector windowRepDefPageData_; + std::unique_ptr<::arrow::util::RleDecoder> windowRepeatDecoder_; + std::unique_ptr<::arrow::util::RleDecoder> windowDefineDecoder_; + int32_t windowRepDefsRemaining_{0}; int32_t repDefMemoryLimit_{16L << 20}; int64_t totalRefDefBytes_{0}; int32_t decodeRepDefPageCount_{10}; From 7389117a377a91d7b43f7a7747a73d9539895191 Mon Sep 17 00:00:00 2001 From: "huangjianfeng.118" Date: Thu, 23 Jul 2026 19:38:26 +0800 Subject: [PATCH 3/3] Optimize Parquet lazy repdef window decode --- bolt/dwio/parquet/reader/PageReader.cpp | 376 ++++++++++++++++------- bolt/dwio/parquet/reader/PageReader.h | 15 +- bolt/dwio/parquet/reader/ParquetData.cpp | 3 +- 3 files changed, 274 insertions(+), 120 deletions(-) diff --git a/bolt/dwio/parquet/reader/PageReader.cpp b/bolt/dwio/parquet/reader/PageReader.cpp index a061c8c3b..b79002d8c 100644 --- a/bolt/dwio/parquet/reader/PageReader.cpp +++ b/bolt/dwio/parquet/reader/PageReader.cpp @@ -103,7 +103,14 @@ void PageReader::seekToPage(int64_t row, bool keepRepDefRawData) { numRowsInPage_ = 0; break; } - if (row != kRepDefOnly && hasChunkRepDefs_ && !isTopLevel_ && + if (row != kRepDefOnly && windowRepDefMode_ && hasChunkRepDefs_ && + !isTopLevel_ && maxRepeat_ > 0) { + while (pageIndex_ + 1 >= static_cast(numLeavesInPage_.size()) && + (!preloadedRepDefs_.empty() || windowRepDefsRemaining_ > 0)) { + ensureNumLeavesInPage(pageIndex_ + 1); + } + } else if ( + row != kRepDefOnly && hasChunkRepDefs_ && !isTopLevel_ && maxRepeat_ > 0) { // In non-window mode we still need the existing sampled-preload guard // to materialize leaf counts before crossing the sampling boundary. @@ -845,8 +852,13 @@ void PageReader::preloadPageRepDefs(const bool keepRepDefRawData) { void PageReader::preloadRepDefs() { const auto startNs = getCurrentTimeNano(); hasChunkRepDefs_ = true; + windowRepDefMode_ = false; + windowTopLevelOffsets_.clear(); windowRepDefPageData_.clear(); windowRepDefsRemaining_ = 0; + windowCurrentPageLeafCount_ = 0; + windowDecodedPageCount_ = 0; + windowTopLevelOffsetBase_ = 0; bool startWithDictinoaryPage = cryptoCtx_.startDecryptWithDictionaryPage; if (maxRepeat_ > 0 && maxDefine_ > 0) { const int32_t samplePages = decodeRepDefPageCount_; @@ -866,12 +878,12 @@ void PageReader::preloadRepDefs() { auto avgRefDefBytes = totalRefDefBytes_ / samplePages + 1; auto avgPageLen = pageStart_ / samplePages + 1; VLOG(1) << __FUNCTION__ << "[" << this - << "] : def/rep levels = " << definitionLevels_.size() - << ", avgPageLen = " << avgPageLen - << ", chunkSize = " << chunkSize_ - << ", pagesCntPerBatch = " << pageCntPerBatch - << ", samplePages = " << samplePages << ", repDefMemoryLimit_ " - << repDefMemoryLimit_; + << "] : def/rep levels = " << definitionLevels_.size() + << ", avgPageLen = " << avgPageLen + << ", chunkSize = " << chunkSize_ + << ", pagesCntPerBatch = " << pageCntPerBatch + << ", samplePages = " << samplePages << ", repDefMemoryLimit_ " + << repDefMemoryLimit_; pageCnt = 0; do { while (pageStart_ < chunkSize_ && ++pageCnt <= pageCntPerBatch) { @@ -914,6 +926,7 @@ void PageReader::preloadRepDefs() { void PageReader::preloadRepDefsRawOnly() { const auto startNs = getCurrentTimeNano(); hasChunkRepDefs_ = true; + windowRepDefMode_ = true; const bool startWithDictionaryPage = cryptoCtx_.startDecryptWithDictionaryPage; totalRefDefBytes_ = 0; @@ -921,16 +934,17 @@ void PageReader::preloadRepDefsRawOnly() { windowRepeatDecoder_.reset(); windowDefineDecoder_.reset(); windowRepDefsRemaining_ = 0; + windowCurrentPageLeafCount_ = 0; + windowDecodedPageCount_ = 0; windowRepDefPageData_.clear(); + windowTopLevelOffsets_.clear(); + windowTopLevelOffsetBase_ = 0; while (pageStart_ < chunkSize_) { preloadedRepDefs_.emplace_back(raw_vector(&pool_)); preloadPageRepDefs(true); if (preloadedRepDefs_.back().empty()) { preloadedRepDefs_.pop_back(); - } else { - numLeavesInPage_.push_back( - countLeavesInRepDefPage(preloadedRepDefs_.back())); } } @@ -947,8 +961,7 @@ void PageReader::preloadRepDefsRawOnly() { cryptoCtx_.startDecryptWithDictionaryPage = startWithDictionaryPage; addPageReaderStat("pageReaderPreloadRepDefsRawOnlyCount", 1); addPageReaderNanos( - "pageReaderPreloadRepDefsRawOnlyTimeNs", - getCurrentTimeNano() - startNs); + "pageReaderPreloadRepDefsRawOnlyTimeNs", getCurrentTimeNano() - startNs); } int32_t PageReader::countLeavesInRepDefPage( @@ -971,32 +984,52 @@ int32_t PageReader::countLeavesInRepDefPage( const auto defineLength = readField(rawData); BOLT_CHECK_LE(rawData + defineLength, pageEnd); - auto defineDecoder = ::arrow::util::RleDecoder( - reinterpret_cast(rawData), - defineLength, - ::arrow::bit_util::NumRequiredBits(maxDefine_)); - - constexpr int32_t kMaxRepDefDecodeBatch = 64 * 1024; - raw_vector definitionScratch(&pool_); - raw_vector nullScratch(&pool_); - definitionScratch.resize(std::min(numRepDefsInPage, kMaxRepDefDecodeBatch)); + if (leafInfo_.rep_level == 0) { + return numRepDefsInPage; + } - int32_t decoded = 0; + auto bitReader = ::arrow::bit_util::BitReader( + reinterpret_cast(rawData), defineLength); + const auto bitWidth = ::arrow::bit_util::NumRequiredBits(maxDefine_); + const auto presentLevel = leafInfo_.repeated_ancestor_def_level; + auto remaining = numRepDefsInPage; int32_t numLeaves = 0; - while (decoded < numRepDefsInPage) { - const auto batchSize = - std::min(numRepDefsInPage - decoded, kMaxRepDefDecodeBatch); - defineDecoder.GetBatch(definitionScratch.data(), batchSize); - nullScratch.resize(bits::nwords(batchSize)); - arrow::ValidityBitmapInputOutput bits; - bits.values_read_upper_bound = batchSize; - bits.values_read = 0; - bits.null_count = 0; - bits.valid_bits = reinterpret_cast(nullScratch.data()); - bits.valid_bits_offset = 0; - DefLevelsToBitmap(definitionScratch.data(), batchSize, leafInfo_, &bits); - numLeaves += bits.values_read; - decoded += batchSize; + + while (remaining > 0) { + uint32_t indicator = 0; + BOLT_CHECK(bitReader.GetVlqInt(&indicator), "Invalid definition RLE run"); + const bool isLiteral = indicator & 1; + const uint32_t count = indicator >> 1; + + if (isLiteral) { + BOLT_CHECK_GT(count, 0); + BOLT_CHECK_LE(count, static_cast(INT32_MAX) / 8); + const auto literalCount = + std::min(remaining, static_cast(count * 8)); + for (auto i = 0; i < literalCount; ++i) { + uint64_t value = 0; + BOLT_CHECK( + bitReader.GetValue(bitWidth, &value), + "Invalid definition literal run"); + numLeaves += value >= presentLevel ? 1 : 0; + } + remaining -= literalCount; + } else { + BOLT_CHECK_GT(count, 0); + BOLT_CHECK_LE(count, static_cast(INT32_MAX)); + uint64_t value = 0; + BOLT_CHECK( + bitReader.GetAligned( + static_cast(::arrow::bit_util::CeilDiv(bitWidth, 8)), + &value), + "Invalid definition repeated run"); + const auto repeatCount = + std::min(remaining, static_cast(count)); + if (value >= presentLevel) { + numLeaves += repeatCount; + } + remaining -= repeatCount; + } } return numLeaves; } @@ -1072,15 +1105,146 @@ void PageReader::compactWindowRepDefs() { } definitionLevels_.resize(keptLevels); repetitionLevels_.resize(keptLevels); + compactWindowTopLevelOffsets(repDefEnd_); repDefBegin_ = 0; repDefEnd_ = 0; } +void PageReader::compactWindowTopLevelOffsets(int32_t consumedLevels) { + if (consumedLevels == 0 || windowTopLevelOffsets_.empty()) { + return; + } + + windowTopLevelOffsetBase_ += consumedLevels; + auto firstKept = std::lower_bound( + windowTopLevelOffsets_.begin(), + windowTopLevelOffsets_.end(), + windowTopLevelOffsetBase_); + const auto kept = windowTopLevelOffsets_.end() - firstKept; + if (kept > 0 && firstKept != windowTopLevelOffsets_.begin()) { + memmove( + windowTopLevelOffsets_.data(), + &*firstKept, + kept * sizeof(windowTopLevelOffsets_[0])); + } + windowTopLevelOffsets_.resize(kept); +} + +int32_t PageReader::appendWindowTopLevelOffsets(int32_t begin, int32_t end) { + if (maxRepeat_ == 0) { + return std::max(end - begin, 0); + } + + int32_t added = 0; + for (int32_t i = begin; i < end; ++i) { + if (repetitionLevels_[i] == 0) { + windowTopLevelOffsets_.push_back(windowTopLevelOffsetBase_ + i); + ++added; + } + } + return added; +} + +int32_t PageReader::decodeWindowRepDefsBatch(int32_t batchSize) { + if (batchSize <= 0) { + return 0; + } + + const auto begin = static_cast(definitionLevels_.size()); + const auto end = begin + batchSize; + repetitionLevels_.resize(end); + BOLT_CHECK_NOT_NULL( + windowRepeatDecoder_, "Missing incremental repetition decoder"); + windowRepeatDecoder_->GetBatch(repetitionLevels_.data() + begin, batchSize); + + definitionLevels_.resize(end); + BOLT_CHECK_NOT_NULL( + windowDefineDecoder_, "Missing incremental definition decoder"); + windowDefineDecoder_->GetBatch(definitionLevels_.data() + begin, batchSize); + + windowRepDefsRemaining_ -= batchSize; + leafNulls_.resize(bits::nwords(leafNullsSize_ + batchSize)); + const auto numLeaves = getLengthsAndNulls( + LevelMode::kNulls, + leafInfo_, + begin, + end, + batchSize, + nullptr, + leafNulls_.data(), + leafNullsSize_); + leafNullsSize_ += numLeaves; + windowCurrentPageLeafCount_ += numLeaves; + appendWindowTopLevelOffsets(begin, end); + finishCurrentWindowRepDefPage(); + return numLeaves; +} + +void PageReader::finishCurrentWindowRepDefPage() { + if (windowRepDefsRemaining_ == 0 && !windowRepDefPageData_.empty()) { + if (windowDecodedPageCount_ == numLeavesInPage_.size()) { + numLeavesInPage_.push_back(windowCurrentPageLeafCount_); + } else { + BOLT_CHECK_LT(windowDecodedPageCount_, numLeavesInPage_.size()); + BOLT_DCHECK_EQ( + numLeavesInPage_[windowDecodedPageCount_], + windowCurrentPageLeafCount_); + } + ++windowDecodedPageCount_; + windowCurrentPageLeafCount_ = 0; + windowRepDefPageData_.clear(); + } +} + +void PageReader::ensureNumLeavesInPage(int32_t targetPageIndex) { + if (windowRepDefMode_) { + if (targetPageIndex < static_cast(numLeavesInPage_.size())) { + return; + } + + auto nextPageToCount = static_cast(numLeavesInPage_.size()); + if (!windowRepDefPageData_.empty() && + windowDecodedPageCount_ == nextPageToCount) { + numLeavesInPage_.push_back( + countLeavesInRepDefPage(windowRepDefPageData_)); + nextPageToCount = static_cast(numLeavesInPage_.size()); + if (targetPageIndex < static_cast(numLeavesInPage_.size())) { + return; + } + } + + const auto firstPreloadedPageIndex = + windowDecodedPageCount_ + (windowRepDefPageData_.empty() ? 0 : 1); + auto pagesToSkip = nextPageToCount - firstPreloadedPageIndex; + BOLT_CHECK_GE(pagesToSkip, 0); + for (const auto& repDefData : preloadedRepDefs_) { + if (repDefData.empty()) { + continue; + } + if (pagesToSkip > 0) { + --pagesToSkip; + continue; + } + numLeavesInPage_.push_back(countLeavesInRepDefPage(repDefData)); + if (targetPageIndex < static_cast(numLeavesInPage_.size())) { + break; + } + } + return; + } + + while (targetPageIndex >= static_cast(numLeavesInPage_.size()) && + !preloadedRepDefs_.empty()) { + loadMoreRepDefs(); + } +} + bool PageReader::loadNextWindowRepDefPage() { if (windowRepDefsRemaining_ > 0) { return true; } + finishCurrentWindowRepDefPage(); windowRepeatDecoder_.reset(); windowDefineDecoder_.reset(); while (!preloadedRepDefs_.empty()) { @@ -1095,8 +1259,10 @@ bool PageReader::loadNextWindowRepDefPage() { windowRepDefPageData_.data() + windowRepDefPageData_.size(); windowRepDefsRemaining_ = readField(rawData); if (windowRepDefsRemaining_ <= 0) { + finishCurrentWindowRepDefPage(); continue; } + windowCurrentPageLeafCount_ = 0; const auto repeatLength = readField(rawData); BOLT_CHECK_GE(repeatLength, 0); @@ -1118,17 +1284,6 @@ bool PageReader::loadNextWindowRepDefPage() { return false; } -int32_t PageReader::countTopLevelRepDefs(int32_t begin, int32_t end) const { - if (maxRepeat_ == 0) { - return std::max(end - begin, 0); - } - int32_t topLevelRows = 0; - for (int32_t i = begin; i < end; ++i) { - topLevelRows += repetitionLevels_[i] == 0 ? 1 : 0; - } - return topLevelRows; -} - void PageReader::loadWindowRepDefs(int32_t targetTopLevelRows) { const auto startNs = getCurrentTimeNano(); if (targetTopLevelRows <= 0) { @@ -1137,8 +1292,8 @@ void PageReader::loadWindowRepDefs(int32_t targetTopLevelRows) { constexpr int32_t kMaxRepDefDecodeBatch = 64 * 1024; compactConsumedLeafNulls(); - auto topLevelRows = - countTopLevelRepDefs(repDefBegin_, definitionLevels_.size()); + compactWindowTopLevelOffsets(repDefBegin_); + auto topLevelRows = static_cast(windowTopLevelOffsets_.size()); int32_t decodeBatchCount = 0; while (topLevelRows < targetTopLevelRows && loadNextWindowRepDefPage()) { const auto batchSize = @@ -1146,40 +1301,17 @@ void PageReader::loadWindowRepDefs(int32_t targetTopLevelRows) { if (batchSize <= 0) { continue; } - - const auto begin = static_cast(definitionLevels_.size()); - const auto end = begin + batchSize; - repetitionLevels_.resize(end); - BOLT_CHECK_NOT_NULL( - windowRepeatDecoder_, "Missing incremental repetition decoder"); - windowRepeatDecoder_->GetBatch(repetitionLevels_.data() + begin, batchSize); - - definitionLevels_.resize(end); - BOLT_CHECK_NOT_NULL( - windowDefineDecoder_, "Missing incremental definition decoder"); - windowDefineDecoder_->GetBatch(definitionLevels_.data() + begin, batchSize); - - windowRepDefsRemaining_ -= batchSize; - leafNulls_.resize(bits::nwords(leafNullsSize_ + batchSize)); - const auto numLeaves = getLengthsAndNulls( - LevelMode::kNulls, - leafInfo_, - begin, - end, - batchSize, - nullptr, - leafNulls_.data(), - leafNullsSize_); - leafNullsSize_ += numLeaves; - topLevelRows += countTopLevelRepDefs(begin, end); + const auto topLevelCountBefore = + static_cast(windowTopLevelOffsets_.size()); + decodeWindowRepDefsBatch(batchSize); + topLevelRows += static_cast(windowTopLevelOffsets_.size()) - + topLevelCountBefore; ++decodeBatchCount; } addPageReaderStat("pageReaderLoadWindowRepDefsCount", 1); - addPageReaderStat( - "pageReaderLoadWindowRepDefsBatches", decodeBatchCount); + addPageReaderStat("pageReaderLoadWindowRepDefsBatches", decodeBatchCount); addPageReaderNanos( - "pageReaderLoadWindowRepDefsTimeNs", - getCurrentTimeNano() - startNs); + "pageReaderLoadWindowRepDefsTimeNs", getCurrentTimeNano() - startNs); } void PageReader::decodeRepDefs(int32_t numTopLevelRows) { @@ -1222,51 +1354,62 @@ void PageReader::decodeRepDefs(int32_t numTopLevelRows) { } int32_t numLevels = definitionLevels_.size(); int32_t topFound = 0; - int32_t i = repDefBegin_; int32_t loadMoreCount = 0; - auto foundTopLevel = [&]() { - const auto scanBegin = i; - for (; i < numLevels; ++i) { - if (repetitionLevels_[i] == 0) { - ++topFound; - if (topFound == numTopLevelRows + 1) { - break; + if (maxRepeat_ > 0) { + if (windowPreload) { + const auto findStartNs = getCurrentTimeNano(); + compactWindowTopLevelOffsets(repDefBegin_); + topFound = std::min( + numTopLevelRows + 1, + static_cast(windowTopLevelOffsets_.size())); + if (topFound == numTopLevelRows + 1) { + repDefEnd_ = + windowTopLevelOffsets_[numTopLevelRows] - windowTopLevelOffsetBase_; + } else { + repDefEnd_ = numLevels; + } + findTopLevelLevels = + repDefEnd_ > repDefBegin_ ? repDefEnd_ - repDefBegin_ : 0; + findTopLevelTimeNs += getCurrentTimeNano() - findStartNs; + } else { + int32_t i = repDefBegin_; + auto foundTopLevel = [&]() { + const auto scanBegin = i; + for (; i < numLevels; ++i) { + if (repetitionLevels_[i] == 0) { + ++topFound; + if (topFound == numTopLevelRows + 1) { + break; + } + } + } + auto scannedLevels = i - scanBegin; + if (i < numLevels && topFound == numTopLevelRows + 1) { + ++scannedLevels; } + repDefEnd_ = i; + return scannedLevels; + }; + const auto findStartNs = getCurrentTimeNano(); + findTopLevelLevels += foundTopLevel(); + findTopLevelTimeNs += getCurrentTimeNano() - findStartNs; + while (repDefEnd_ == numLevels && topFound < numTopLevelRows + 1 && + (!preloadedRepDefs_.empty() || windowRepDefsRemaining_ > 0)) { + const auto loadMoreStartNs = getCurrentTimeNano(); + loadMoreRepDefs(); + ++loadMoreCount; + numLevels = definitionLevels_.size(); + i = repDefEnd_; + findTopLevelLevels += foundTopLevel(); + loadMoreTimeNs += getCurrentTimeNano() - loadMoreStartNs; + BOLT_CHECK(topFound == numTopLevelRows + 1 || repDefEnd_ == numLevels); } } - auto scannedLevels = i - scanBegin; - if (i < numLevels && topFound == numTopLevelRows + 1) { - ++scannedLevels; - } - repDefEnd_ = i; - return scannedLevels; - }; - - if (maxRepeat_ > 0) { - const auto findStartNs = getCurrentTimeNano(); - findTopLevelLevels += foundTopLevel(); - findTopLevelTimeNs += getCurrentTimeNano() - findStartNs; } else { - repDefEnd_ = i + numTopLevelRows; + repDefEnd_ = repDefBegin_ + numTopLevelRows; } if (maxRepeat_ > 0 && maxDefine_ > 0) { - // definitionLevels_ has been consumed, decode more if any - while (repDefEnd_ == numLevels && topFound < numTopLevelRows + 1 && - (!preloadedRepDefs_.empty() || windowRepDefsRemaining_ > 0)) { - const auto loadMoreStartNs = getCurrentTimeNano(); - if (windowPreload) { - loadWindowRepDefs(targetTopLevelRows); - } else { - loadMoreRepDefs(); - } - ++loadMoreCount; - numLevels = definitionLevels_.size(); - i = repDefEnd_; - findTopLevelLevels += foundTopLevel(); - loadMoreTimeNs += getCurrentTimeNano() - loadMoreStartNs; - BOLT_CHECK(topFound == numTopLevelRows + 1 || repDefEnd_ == numLevels); - } // after topFound done, left decoded rep/def is less than 1 page, decode // more if (!windowPreload && !numLeavesInPage_.empty() && @@ -1441,6 +1584,7 @@ void PageReader::decodeRepDefsFromBuffer() { numLeavesInPage_.push_back(numLeaves); } preloadedRepDefs_.pop_front(); + windowRepDefMode_ = false; } int32_t PageReader::getLengthsAndNulls( diff --git a/bolt/dwio/parquet/reader/PageReader.h b/bolt/dwio/parquet/reader/PageReader.h index 842b4f953..578ac8875 100644 --- a/bolt/dwio/parquet/reader/PageReader.h +++ b/bolt/dwio/parquet/reader/PageReader.h @@ -94,6 +94,7 @@ class PageReader { repetitionLevels_(&pool_), nullConcatenation_(pool_), windowRepDefPageData_(&pool_), + windowTopLevelOffsets_(&pool_), statis_(statis) { type_->makeLevelInfo(leafInfo_); pageOrdinal_ = 0; @@ -115,7 +116,8 @@ class PageReader { definitionLevels_(&pool_), repetitionLevels_(&pool_), nullConcatenation_(pool_), - windowRepDefPageData_(&pool_) { + windowRepDefPageData_(&pool_), + windowTopLevelOffsets_(&pool_) { pageOrdinal_ = 0; } @@ -479,9 +481,13 @@ class PageReader { int32_t countLeavesInRepDefPage(const raw_vector& repDefData); void compactConsumedLeafNulls(); void compactWindowRepDefs(); + void compactWindowTopLevelOffsets(int32_t consumedLevels); + int32_t appendWindowTopLevelOffsets(int32_t begin, int32_t end); + int32_t decodeWindowRepDefsBatch(int32_t batchSize); + void finishCurrentWindowRepDefPage(); + void ensureNumLeavesInPage(int32_t pageIndex); void loadWindowRepDefs(int32_t targetTopLevelRows); bool loadNextWindowRepDefPage(); - int32_t countTopLevelRepDefs(int32_t begin, int32_t end) const; void decodeRepDefsFromBuffer(); memory::MemoryPool& pool_; @@ -639,10 +645,15 @@ class PageReader { // preload undecoded RepDefs std::list> preloadedRepDefs_; + bool windowRepDefMode_{false}; raw_vector windowRepDefPageData_; + raw_vector windowTopLevelOffsets_; + int32_t windowTopLevelOffsetBase_{0}; std::unique_ptr<::arrow::util::RleDecoder> windowRepeatDecoder_; std::unique_ptr<::arrow::util::RleDecoder> windowDefineDecoder_; int32_t windowRepDefsRemaining_{0}; + int32_t windowCurrentPageLeafCount_{0}; + int32_t windowDecodedPageCount_{0}; int32_t repDefMemoryLimit_{16L << 20}; int64_t totalRefDefBytes_{0}; int32_t decodeRepDefPageCount_{10}; diff --git a/bolt/dwio/parquet/reader/ParquetData.cpp b/bolt/dwio/parquet/reader/ParquetData.cpp index bf8e16564..a35da1fa1 100644 --- a/bolt/dwio/parquet/reader/ParquetData.cpp +++ b/bolt/dwio/parquet/reader/ParquetData.cpp @@ -374,8 +374,7 @@ dwio::common::PositionProvider ParquetData::seekToRowGroup(int64_t index) { metadata.total_compressed_size, statis_); reader_->setDecodeRepDefPageCount(decodeRepDefPageCount_); - reader_->setParquetRepDefPreloadWindowCount( - parquetRepDefPreloadWindowCount_); + reader_->setParquetRepDefPreloadWindowCount(parquetRepDefPreloadWindowCount_); reader_->setParquetRepDefMemoryLimit(parquetRepDefMemoryLimit_); if (columnChunkMeta.__isset.crypto_metadata) {