From 121bf9dd8225c64844a163e2e6991c92479f1990 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 01:15:18 +0000 Subject: [PATCH 01/15] Add hybrid scan multifile reader basics --- cpp/CMakeLists.txt | 1 + .../io/experimental/hybrid_scan_multifile.hpp | 146 +++++++++++++++ .../io/parquet/experimental/hybrid_scan.cpp | 15 +- .../experimental/hybrid_scan_helpers.cpp | 166 ++++++++++++------ .../experimental/hybrid_scan_helpers.hpp | 39 ++-- .../parquet/experimental/hybrid_scan_impl.cpp | 25 ++- .../parquet/experimental/hybrid_scan_impl.hpp | 35 ++-- .../experimental/hybrid_scan_multifile.cpp | 60 +++++++ 8 files changed, 381 insertions(+), 106 deletions(-) create mode 100644 cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp create mode 100644 cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp diff --git a/cpp/CMakeLists.txt b/cpp/CMakeLists.txt index 63c292646773..d3576d2a6b5a 100644 --- a/cpp/CMakeLists.txt +++ b/cpp/CMakeLists.txt @@ -627,6 +627,7 @@ add_library( src/io/parquet/experimental/hybrid_scan_chunking.cu src/io/parquet/experimental/hybrid_scan_helpers.cpp src/io/parquet/experimental/hybrid_scan_impl.cpp + src/io/parquet/experimental/hybrid_scan_multifile.cpp src/io/parquet/experimental/hybrid_scan_preprocess.cu src/io/parquet/experimental/page_index_filter.cu src/io/parquet/experimental/page_index_filter_utils.cu diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp new file mode 100644 index 000000000000..13bbf97ee499 --- /dev/null +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -0,0 +1,146 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-License-Identifier: Apache-2.0 + */ + +#pragma once + +#include +#include +#include +#include +#include +#include +#include + +#include +#include + +#include +#include + +namespace CUDF_EXPORT cudf { +namespace io::parquet::experimental::detail { +/** + * @brief Internal experimental Parquet reader optimized for highly selective filters, called a + * Hybrid Scan operation. + */ +class hybrid_scan_reader_impl; +} // namespace io::parquet::experimental::detail +} // namespace CUDF_EXPORT cudf + +//! Using `byte_range_info` from cudf::io::text +using cudf::io::text::byte_range_info; + +namespace CUDF_EXPORT cudf { +namespace io::parquet::experimental { +/** + * @addtogroup io_readers + * @{ + * @file + */ + +/** + * @brief Multi-file variant of the experimental Hybrid Scan Parquet reader + * + * Vectorizes `hybrid_scan_reader` APIs to support multiple Parquet sources. Inputs and outputs are + * indexed by source order except for the row mask which is a single BOOL8 column spanning all rows + * from all sources concatenated in source order, then row-group order within a source. + * + * @note Detailed usage documentation will be added once all APIs are in place. + */ +class hybrid_scan_multifile { + public: + /** + * @brief Constructor for the multi-file experimental Parquet reader + * + * @param footer_bytes Host span of Parquet file footer byte spans, one per source + * @param options Parquet reader options + */ + explicit hybrid_scan_multifile(cudf::host_span const> footer_bytes, + parquet_reader_options const& options); + + /** + * @brief Constructor for the multi-file experimental Parquet reader + * + * @param parquet_metadata Host span of pre-populated Parquet file metadata, one per source + * @param options Parquet reader options + */ + explicit hybrid_scan_multifile(cudf::host_span parquet_metadata, + parquet_reader_options const& options); + + /** + * @brief Destructor for the multi-file experimental Parquet reader + */ + ~hybrid_scan_multifile(); + + /** + * @brief Get the per-source Parquet file footer metadata + * + * @return Vector of file metadata, one per source + */ + [[nodiscard]] std::vector parquet_metadata() const; + + /** + * @brief Get the per-source byte range of the page index in each Parquet file + * + * The returned vector always has one entry per source. A source's entry is a + * default-constructed `byte_range_info{}` if that source has no row groups, no columns, or no + * page index offsets. A `CUDF_LOG_WARN` is emitted once if some sources have a page index and + * others do not. + * + * @return Vector of page index byte ranges, one per source + */ + [[nodiscard]] std::vector page_index_byte_range() const; + + /** + * @brief Setup the per-source page index within each Parquet file metadata + * + * Materializes `ColumnIndex` and `OffsetIndex` (page index) inside each source's + * `FileMetaData`. The input span size must equal the number of sources. A per-source empty + * span is skipped with a one-time warning. Sources whose corresponding span is non-empty must + * have row groups and valid page index offsets. + * + * @param page_index_bytes Host span of Parquet page index buffer bytes, one per source + */ + void setup_page_index( + cudf::host_span const> page_index_bytes) const; + + /** + * @brief Get all available per-source row group indices from the parquet files + * + * If `options.get_row_groups()` is non-empty, its size must equal the number of sources and it + * is returned as-is. Otherwise builds `[0 .. per_source_num_row_groups[i])` for each source. + * + * @param options Parquet reader options + * @return Vector of row group indices, one inner vector per source + */ + [[nodiscard]] std::vector> all_row_groups( + parquet_reader_options const& options) const; + + /** + * @brief Get the total number of top-level rows in the per-source row groups + * + * @param row_group_indices Input per-source row group indices (one inner vector per source) + * @return Total number of top-level rows across all sources + */ + [[nodiscard]] size_type total_rows_in_row_groups( + cudf::host_span const> row_group_indices) const; + + /** + * @brief Resets the current column selection + * + * Resets the current column selection state forcing column re-selection in subsequent filter, + * byte range, setup chunking and materialization APIs. This is useful if the filter expression + * has been cascaded (and-ed) to include new columns. + */ + void reset_column_selection() const; + + private: + std::unique_ptr _impl; +}; + +/** @} */ // end of group + +} // namespace io::parquet::experimental +} // namespace CUDF_EXPORT cudf diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index b6243fafe60a..3223ca99c8f8 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -15,33 +15,36 @@ namespace cudf::io::parquet::experimental { hybrid_scan_reader::hybrid_scan_reader(cudf::host_span footer_bytes, parquet_reader_options const& options) - : _impl{std::make_unique(footer_bytes, options)} { + auto const footers = std::vector>{footer_bytes}; + _impl = std::make_unique(footers, options); } hybrid_scan_reader::hybrid_scan_reader(FileMetaData const& parquet_metadata, parquet_reader_options const& options) - : _impl{std::make_unique(parquet_metadata, options)} { + auto const metadatas = std::vector{parquet_metadata}; + _impl = std::make_unique(metadatas, options); } hybrid_scan_reader::~hybrid_scan_reader() = default; [[nodiscard]] text::byte_range_info hybrid_scan_reader::page_index_byte_range() const { - return _impl->page_index_byte_range(); + return _impl->page_index_byte_range().front(); } [[nodiscard]] FileMetaData hybrid_scan_reader::parquet_metadata() const { - return _impl->parquet_metadata(); + return _impl->parquet_metadata().front(); } void hybrid_scan_reader::setup_page_index(cudf::host_span page_index_bytes) const { CUDF_FUNC_RANGE(); - return _impl->setup_page_index(page_index_bytes); + auto const per_source = std::vector>{page_index_bytes}; + return _impl->setup_page_index(per_source); } std::vector hybrid_scan_reader::all_row_groups( @@ -53,7 +56,7 @@ std::vector hybrid_scan_reader::all_row_groups( // If row groups are specified in parquet reader options, return them as is if (options.get_row_groups().size() == 1) { return options.get_row_groups().front(); } - return _impl->all_row_groups(options); + return _impl->all_row_groups(options).front(); } size_type hybrid_scan_reader::total_rows_in_row_groups( diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index f6e6985ea4ea..e93054073637 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -72,23 +72,34 @@ metadata::metadata(cudf::host_span footer_bytes) sanitize_schema(); } -aggregate_reader_metadata::aggregate_reader_metadata(FileMetaData const& parquet_metadata, - bool use_arrow_schema, - bool has_cols_from_mismatched_srcs) +aggregate_reader_metadata::aggregate_reader_metadata( + cudf::host_span const> footer_bytes, + bool use_arrow_schema, + bool has_cols_from_mismatched_srcs) : aggregate_reader_metadata_base(host_span const>{}, false, false) { - // Just copy over the FileMetaData struct to the internal metadata struct - per_file_metadata.emplace_back(metadata{parquet_metadata}); + CUDF_EXPECTS(not footer_bytes.empty(), "At least one source must be provided"); + per_file_metadata.reserve(footer_bytes.size()); + std::transform(footer_bytes.begin(), + footer_bytes.end(), + std::back_inserter(per_file_metadata), + [](auto const& fb) { return metadata{fb}; }); initialize_internals(use_arrow_schema, has_cols_from_mismatched_srcs); } -aggregate_reader_metadata::aggregate_reader_metadata(cudf::host_span footer_bytes, - bool use_arrow_schema, - bool has_cols_from_mismatched_srcs) +aggregate_reader_metadata::aggregate_reader_metadata( + cudf::host_span parquet_metadatas, + bool use_arrow_schema, + bool has_cols_from_mismatched_srcs) : aggregate_reader_metadata_base(host_span const>{}, false, false) { - // Re-initialize internal variables here as base class was initialized without a source - per_file_metadata.emplace_back(metadata{footer_bytes}); + CUDF_EXPECTS(not parquet_metadatas.empty(), "At least one source must be provided"); + per_file_metadata.reserve(parquet_metadatas.size()); + // Just copy over the FileMetaData structs to the internal metadata structs + std::transform(parquet_metadatas.begin(), + parquet_metadatas.end(), + std::back_inserter(per_file_metadata), + [](auto const& parquet_metadata) { return metadata{parquet_metadata}; }); initialize_internals(use_arrow_schema, has_cols_from_mismatched_srcs); } @@ -102,13 +113,15 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, // Force all non-nullable (REQUIRED) columns to be nullable without modifying REPEATED columns to // preserve list structures - auto& schema = per_file_metadata.front().schema; - std::for_each(schema.begin() + 1, schema.end(), [](auto& col) { - // TODO: Store information of whichever column schema we modified here and restore it to - // `REQUIRED` if we end up not pruning any pages out of it - if (col.repetition_type == FieldRepetitionType::REQUIRED) { - col.repetition_type = FieldRepetitionType::OPTIONAL; - } + std::for_each(per_file_metadata.begin(), per_file_metadata.end(), [](auto& pfm) { + auto& schema = pfm.schema; + std::for_each(schema.begin() + 1, schema.end(), [](auto& col) { + // TODO: Store information of whichever column schema we modified here and restore it to + // `REQUIRED` if we end up not pruning any pages out of it + if (col.repetition_type == FieldRepetitionType::REQUIRED) { + col.repetition_type = FieldRepetitionType::OPTIONAL; + } + }); }); // Collect and apply arrow:schema from Parquet's key value metadata section @@ -122,53 +135,90 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, } } -text::byte_range_info aggregate_reader_metadata::page_index_byte_range() const +std::vector aggregate_reader_metadata::page_index_byte_range() const { - auto& schema = per_file_metadata.front(); - auto& row_groups = schema.row_groups; - - if (row_groups.size() and row_groups.front().columns.size()) { - auto const min_offset = schema.row_groups.front().columns.front().column_index_offset; - auto const& last_col = schema.row_groups.back().columns.back(); - auto const max_offset = last_col.offset_index_offset + last_col.offset_index_length; - return {min_offset, (max_offset - min_offset)}; - } + std::vector page_index_byte_ranges; + page_index_byte_ranges.reserve(per_file_metadata.size()); + std::transform(per_file_metadata.begin(), + per_file_metadata.end(), + std::back_inserter(page_index_byte_ranges), + [](auto const& pfm) -> text::byte_range_info { + auto const& row_groups = pfm.row_groups; + if (row_groups.empty() or row_groups.front().columns.empty()) { return {}; } + + auto const min_offset = row_groups.front().columns.front().column_index_offset; + auto const& last_col = row_groups.back().columns.back(); + auto const max_offset = + last_col.offset_index_offset + last_col.offset_index_length; + + if (max_offset <= min_offset) { return {}; } + return {min_offset, max_offset - min_offset}; + }); - return {}; + return page_index_byte_ranges; } -FileMetaData aggregate_reader_metadata::parquet_metadata() const +std::vector aggregate_reader_metadata::parquet_metadata() const { - return per_file_metadata.front(); + return {per_file_metadata.begin(), per_file_metadata.end()}; } -void aggregate_reader_metadata::setup_page_index(cudf::host_span page_index_bytes) +void aggregate_reader_metadata::setup_page_index( + cudf::host_span const> page_index_bytes) { - // Return early if empty page index buffer span - if (page_index_bytes.empty()) { - CUDF_LOG_WARN("Hybrid scan reader encountered empty page index buffer"); - return; - } + CUDF_EXPECTS(page_index_bytes.size() == per_file_metadata.size(), + "Page index byte span count must equal the number of sources"); + + auto iter = cuda::zip_iterator(page_index_bytes.begin(), per_file_metadata.begin()); + std::for_each(iter, iter + page_index_bytes.size(), [&](auto const& pair) { + auto const& pgidx_bytes = cuda::std::get<0>(pair); + + // Return early if empty page index buffer span + if (pgidx_bytes.empty()) { return; } - // Get the file metadata and setup the page index - auto& file_metadata = per_file_metadata.front(); - auto const& row_groups = file_metadata.row_groups; + // Get the file metadata and setup the page index + auto& file_metadata = cuda::std::get<1>(pair); + auto const& row_groups = file_metadata.row_groups; - // Check for empty parquet file - CUDF_EXPECTS(not row_groups.empty() and not row_groups.front().columns.empty(), - "No column chunks in Parquet schema to read page index for"); + // Check for empty parquet file + CUDF_EXPECTS(not row_groups.empty() and not row_groups.front().columns.empty(), + "No column chunks in Parquet schema to read page index for"); - // Set the first ColumnChunk's offset of ColumnIndex as the adjusted zero offset - int64_t const min_offset = row_groups.front().columns.front().column_index_offset; + // Set the first ColumnChunk's offset of ColumnIndex as the adjusted zero offset + int64_t const min_offset = row_groups.front().columns.front().column_index_offset; - // Check if the page index buffer is valid - { - auto const& last_col = row_groups.back().columns.back(); - auto const max_offset = last_col.offset_index_offset + last_col.offset_index_length; - CUDF_EXPECTS(max_offset > min_offset, "Encountered an invalid page index buffer"); + // Check if the page index buffer is valid + { + auto const& last_col = row_groups.back().columns.back(); + auto const max_offset = last_col.offset_index_offset + last_col.offset_index_length; + CUDF_EXPECTS(max_offset > min_offset, "Encountered an invalid page index buffer"); + } + + file_metadata.setup_page_index(pgidx_bytes, min_offset); + }); +} + +std::vector> aggregate_reader_metadata::all_row_groups( + parquet_reader_options const& options) const +{ + auto const& opts_row_groups = options.get_row_groups(); + if (not opts_row_groups.empty()) { + CUDF_EXPECTS(opts_row_groups.size() == per_file_metadata.size(), + "Row groups in parquet reader options must have one inner vector per source"); + return opts_row_groups; } - file_metadata.setup_page_index(page_index_bytes, min_offset); + std::vector> row_groups; + row_groups.reserve(per_file_metadata.size()); + std::transform(per_file_metadata.begin(), + per_file_metadata.end(), + std::back_inserter(row_groups), + [](auto const& pfm) { + std::vector indices(pfm.row_groups.size()); + std::iota(indices.begin(), indices.end(), size_type{0}); + return indices; + }); + return row_groups; } size_type aggregate_reader_metadata::total_rows_in_row_groups( @@ -231,8 +281,8 @@ aggregate_reader_metadata::select_payload_columns( return filter_columns_set; }; - // If payload columns are specified, only select payload columns that do not appear in the filter - // expression + // If payload columns are specified, only select payload columns that do not appear in the + // filter expression if (payload_column_names.has_value()) { valid_payload_columns = *payload_column_names; // Remove filter columns from the provided payload column names @@ -634,12 +684,12 @@ std::reference_wrapper named_to_reference_converter::visi { // Map the column index to its name auto const col_name_iter = _column_indices_to_names.find(expr.get_column_index()); - CUDF_EXPECTS( - col_name_iter != _column_indices_to_names.end(), - "Column index in the filter expression not found in the column indices to names map. Note that " - "only top-level columns except structs and lists are supported in " - "Parquet filter expression", - std::invalid_argument); + CUDF_EXPECTS(col_name_iter != _column_indices_to_names.end(), + "Column index in the filter expression not found in the column indices to names " + "map. Note that " + "only top-level columns except structs and lists are supported in " + "Parquet filter expression", + std::invalid_argument); auto const col_name = col_name_iter->second; auto col_index_it = _column_name_to_index.find(col_name); CUDF_EXPECTS(col_index_it != _column_name_to_index.end(), diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index df17ae493cd2..9a8879163782 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -80,22 +80,22 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { /** * @brief Constructor for aggregate_reader_metadata * - * @param footer_bytes Host span of Parquet file footer buffer bytes + * @param footer_bytes Host span of Parquet file footer buffer bytes, one per source * @param use_arrow_schema Whether to use Arrow schema * @param has_cols_from_mismatched_srcs Whether to have columns from mismatched sources */ - aggregate_reader_metadata(cudf::host_span footer_bytes, + aggregate_reader_metadata(cudf::host_span const> footer_bytes, bool use_arrow_schema, bool has_cols_from_mismatched_srcs); /** * @brief Constructor for aggregate_reader_metadata * - * @param parquet_metadata Pre-populated Parquet file metadata + * @param parquet_metadatas Host span of pre-populated Parquet file metadata, one per source * @param use_arrow_schema Whether to use Arrow schema * @param has_cols_from_mismatched_srcs Whether to have columns from mismatched sources */ - aggregate_reader_metadata(FileMetaData const& parquet_metadata, + aggregate_reader_metadata(cudf::host_span parquet_metadatas, bool use_arrow_schema, bool has_cols_from_mismatched_srcs); @@ -110,21 +110,38 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { void initialize_internals(bool use_arrow_schema, bool has_cols_from_mismatched_srcs); /** - * @brief Fetch the byte range of the page index in the Parquet file + * @brief Fetch the byte range of the page index in each Parquet file + * + * @return Vector of byte ranges of the page index, one per source */ - [[nodiscard]] text::byte_range_info page_index_byte_range() const; + [[nodiscard]] std::vector page_index_byte_range() const; /** - * @brief Get the Parquet file metadata + * @brief Get the Parquet file metadata for every source + * + * @return Vector of file metadata, one per source */ - [[nodiscard]] FileMetaData parquet_metadata() const; + [[nodiscard]] std::vector parquet_metadata() const; /** - * @brief Setup and populate the page index structs in `FileMetaData` + * @brief Setup and populate the page index structs in every source's `FileMetaData` + * + * @param page_index_bytes Host span of Parquet page index buffer bytes, one per source + */ + void setup_page_index(cudf::host_span const> page_index_bytes); + + /** + * @brief Get all available row group indices, one inner vector per source + * + * If `options.get_row_groups()` is non-empty, validates that its size equals the number of + * sources and returns it as-is. Otherwise returns `[0 .. per_source_num_row_groups[i])` for + * each source. * - * @param page_index_bytes Host span of Parquet page index buffer bytes + * @param options Parquet reader options + * @return Vector of row group indices, one inner vector per source */ - void setup_page_index(cudf::host_span page_index_bytes); + [[nodiscard]] std::vector> all_row_groups( + parquet_reader_options const& options) const; /** * @brief Get the total number of top-level rows in the row groups diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 7b53b4345911..5ef78f19299e 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -64,10 +64,10 @@ namespace { } // namespace -hybrid_scan_reader_impl::hybrid_scan_reader_impl(cudf::host_span footer_bytes, - parquet_reader_options const& options) +hybrid_scan_reader_impl::hybrid_scan_reader_impl( + cudf::host_span const> footer_bytes, + parquet_reader_options const& options) { - // Open and parse the source dataset metadata _metadata = std::make_unique( footer_bytes, options.is_enabled_use_arrow_schema(), @@ -76,28 +76,28 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl(cudf::host_span _extended_metadata = static_cast(_metadata.get()); } -hybrid_scan_reader_impl::hybrid_scan_reader_impl(FileMetaData const& parquet_metadata, - parquet_reader_options const& options) +hybrid_scan_reader_impl::hybrid_scan_reader_impl( + cudf::host_span parquet_metadatas, parquet_reader_options const& options) { _metadata = std::make_unique( - parquet_metadata, + parquet_metadatas, options.is_enabled_use_arrow_schema(), options.get_column_names().has_value() and options.is_enabled_allow_mismatched_pq_schemas()); _extended_metadata = static_cast(_metadata.get()); } -FileMetaData hybrid_scan_reader_impl::parquet_metadata() const +std::vector hybrid_scan_reader_impl::parquet_metadata() const { return _extended_metadata->parquet_metadata(); } -byte_range_info hybrid_scan_reader_impl::page_index_byte_range() const +std::vector hybrid_scan_reader_impl::page_index_byte_range() const { return _extended_metadata->page_index_byte_range(); } void hybrid_scan_reader_impl::setup_page_index( - cudf::host_span page_index_bytes) const + cudf::host_span const> page_index_bytes) const { _extended_metadata->setup_page_index(page_index_bytes); } @@ -200,13 +200,10 @@ void hybrid_scan_reader_impl::select_columns(read_columns_mode read_columns_mode [](auto const& buff) { return inline_column_buffer::empty_like(buff); }); } -std::vector hybrid_scan_reader_impl::all_row_groups( +std::vector> hybrid_scan_reader_impl::all_row_groups( parquet_reader_options const& options) const { - auto const num_row_groups = _extended_metadata->get_num_row_groups(); - auto row_groups_indices = std::vector(num_row_groups); - std::iota(row_groups_indices.begin(), row_groups_indices.end(), size_type{0}); - return row_groups_indices; + return _extended_metadata->all_row_groups(options); } size_type hybrid_scan_reader_impl::total_rows_in_row_groups( diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 3f16699a3d18..0e85f066d347 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -39,44 +39,45 @@ using text::byte_range_info; class hybrid_scan_reader_impl : public parquet::detail::reader_impl { public: /** - * @brief Constructor for the experimental parquet reader implementation to optimally read - * Parquet files subject to highly selective filters + * @brief Constructor for the experimental parquet reader implementation * - * @param footer_bytes Host span of parquet file footer bytes + * @param footer_bytes Span of parquet file footer byte spans, one per source * @param options Parquet reader options */ - explicit hybrid_scan_reader_impl(cudf::host_span footer_bytes, - parquet_reader_options const& options); + explicit hybrid_scan_reader_impl( + cudf::host_span const> footer_bytes, + parquet_reader_options const& options); /** - * @brief Constructor for the experimental parquet reader implementation to optimally read - * Parquet files subject to highly selective filters + * @brief Constructor for the experimental parquet reader implementation * - * @param parquet_metadata Pre-populated Parquet file metadata + * @param parquet_metadatas Span of pre-populated Parquet file metadata, one per source * @param options Parquet reader options */ - explicit hybrid_scan_reader_impl(FileMetaData const& parquet_metadata, + explicit hybrid_scan_reader_impl(cudf::host_span parquet_metadatas, parquet_reader_options const& options); /** - * @copydoc cudf::io::experimental::hybrid_scan::parquet_metadata + * @copydoc cudf::io::experimental::hybrid_scan_multifile::parquet_metadata */ - [[nodiscard]] FileMetaData parquet_metadata() const; + [[nodiscard]] std::vector parquet_metadata() const; /** - * @copydoc cudf::io::experimental::hybrid_scan::page_index_byte_range + * @copydoc cudf::io::experimental::hybrid_scan_multifile::page_index_byte_range */ - [[nodiscard]] byte_range_info page_index_byte_range() const; + [[nodiscard]] std::vector page_index_byte_range() const; /** - * @copydoc cudf::io::experimental::hybrid_scan::setup_page_index + * @copydoc cudf::io::experimental::hybrid_scan_multifile::setup_page_index */ - void setup_page_index(cudf::host_span page_index_bytes) const; + void setup_page_index( + cudf::host_span const> page_index_bytes) const; /** - * @copydoc cudf::io::experimental::hybrid_scan::all_row_groups + * @copydoc cudf::io::experimental::hybrid_scan_multifile::all_row_groups */ - [[nodiscard]] std::vector all_row_groups(parquet_reader_options const& options) const; + [[nodiscard]] std::vector> all_row_groups( + parquet_reader_options const& options) const; /** * @copydoc cudf::io::experimental::hybrid_scan::total_rows_in_row_groups diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp new file mode 100644 index 000000000000..32819903f530 --- /dev/null +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -0,0 +1,60 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-License-Identifier: Apache-2.0 + */ + +#include "hybrid_scan_impl.hpp" + +#include +#include + +namespace cudf::io::parquet::experimental { + +hybrid_scan_multifile::hybrid_scan_multifile( + cudf::host_span const> footer_bytes, + parquet_reader_options const& options) + : _impl{std::make_unique(footer_bytes, options)} +{ +} + +hybrid_scan_multifile::hybrid_scan_multifile(cudf::host_span parquet_metadata, + parquet_reader_options const& options) + : _impl{std::make_unique(parquet_metadata, options)} +{ +} + +hybrid_scan_multifile::~hybrid_scan_multifile() = default; + +std::vector hybrid_scan_multifile::parquet_metadata() const +{ + return _impl->parquet_metadata(); +} + +std::vector hybrid_scan_multifile::page_index_byte_range() const +{ + return _impl->page_index_byte_range(); +} + +void hybrid_scan_multifile::setup_page_index( + cudf::host_span const> page_index_bytes) const +{ + CUDF_FUNC_RANGE(); + _impl->setup_page_index(page_index_bytes); +} + +std::vector> hybrid_scan_multifile::all_row_groups( + parquet_reader_options const& options) const +{ + return _impl->all_row_groups(options); +} + +size_type hybrid_scan_multifile::total_rows_in_row_groups( + cudf::host_span const> row_group_indices) const +{ + if (row_group_indices.empty()) { return 0; } + return _impl->total_rows_in_row_groups(row_group_indices); +} + +void hybrid_scan_multifile::reset_column_selection() const { _impl->reset_column_selection(); } + +} // namespace cudf::io::parquet::experimental From b763cdb973c0313beb8e5bb8cdff63bf22de3134 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 01:20:19 +0000 Subject: [PATCH 02/15] Minor --- cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index 13bbf97ee499..6b6418879e21 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -47,7 +47,9 @@ namespace io::parquet::experimental { * indexed by source order except for the row mask which is a single BOOL8 column spanning all rows * from all sources concatenated in source order, then row-group order within a source. * - * @note Detailed usage documentation will be added once all APIs are in place. + * @note Detailed usage documentation will be added once all APIs are in place. This reader will + * eventually move to `hybrid_scan.hpp` and the existing single-file reader (`hybrid_scan_reader`) + * will become its subclass. Only keeping this separate here for now to reduce noise. */ class hybrid_scan_multifile { public: From c9bf41923d20e12e958a7584994279b33acc1fdb Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 01:27:41 +0000 Subject: [PATCH 03/15] Clean up claude's comments --- .../cudf/io/experimental/hybrid_scan_multifile.hpp | 13 ------------- 1 file changed, 13 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index 6b6418879e21..45ea96ab3628 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -86,11 +86,6 @@ class hybrid_scan_multifile { /** * @brief Get the per-source byte range of the page index in each Parquet file * - * The returned vector always has one entry per source. A source's entry is a - * default-constructed `byte_range_info{}` if that source has no row groups, no columns, or no - * page index offsets. A `CUDF_LOG_WARN` is emitted once if some sources have a page index and - * others do not. - * * @return Vector of page index byte ranges, one per source */ [[nodiscard]] std::vector page_index_byte_range() const; @@ -98,11 +93,6 @@ class hybrid_scan_multifile { /** * @brief Setup the per-source page index within each Parquet file metadata * - * Materializes `ColumnIndex` and `OffsetIndex` (page index) inside each source's - * `FileMetaData`. The input span size must equal the number of sources. A per-source empty - * span is skipped with a one-time warning. Sources whose corresponding span is non-empty must - * have row groups and valid page index offsets. - * * @param page_index_bytes Host span of Parquet page index buffer bytes, one per source */ void setup_page_index( @@ -111,9 +101,6 @@ class hybrid_scan_multifile { /** * @brief Get all available per-source row group indices from the parquet files * - * If `options.get_row_groups()` is non-empty, its size must equal the number of sources and it - * is returned as-is. Otherwise builds `[0 .. per_source_num_row_groups[i])` for each source. - * * @param options Parquet reader options * @return Vector of row group indices, one inner vector per source */ From 795f05849d2c677eb057a0856ac6bf8b379a88ab Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 19:48:29 +0000 Subject: [PATCH 04/15] Add gtests --- cpp/tests/CMakeLists.txt | 1 + .../hybrid_scan_multifile_filters_test.cpp | 221 ++++++++++++++++++ 2 files changed, 222 insertions(+) create mode 100644 cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp diff --git a/cpp/tests/CMakeLists.txt b/cpp/tests/CMakeLists.txt index ce7bfdbed0f0..e709a52ace20 100644 --- a/cpp/tests/CMakeLists.txt +++ b/cpp/tests/CMakeLists.txt @@ -349,6 +349,7 @@ ConfigureTest( HYBRID_SCAN_TEST io/experimental/hybrid_scan_composer.cpp io/experimental/hybrid_scan_filters_test.cpp + io/experimental/hybrid_scan_multifile_filters_test.cpp io/experimental/hybrid_scan_test.cpp io/parquet_common.cpp io/parquet_test.cpp diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp new file mode 100644 index 000000000000..0c62cc8b223c --- /dev/null +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -0,0 +1,221 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-License-Identifier: Apache-2.0 + */ + +#include "hybrid_scan_common.hpp" + +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include + +namespace { + +/** + * @brief Struct to hold multifile datasources, and footer buffers along with their byte spans + */ +struct multifile_inputs { + std::vector> datasources; + std::vector> footer_buffers; + std::vector> footer_byte_spans; +}; + +template +multifile_inputs build_multifile_inputs(Buffers const& file_buffers) +{ + multifile_inputs out; + out.datasources.reserve(file_buffers.size()); + out.footer_buffers.reserve(file_buffers.size()); + out.footer_byte_spans.reserve(file_buffers.size()); + for (auto const& buf : file_buffers) { + out.datasources.emplace_back(cudf::io::datasource::create(cudf::host_span( + reinterpret_cast(buf.data()), buf.size()))); + out.footer_buffers.emplace_back( + cudf::io::parquet::fetch_footer_to_host(*out.datasources.back())); + out.footer_byte_spans.emplace_back(*out.footer_buffers.back()); + } + return out; +} + +/** + * @brief Creates a parquet buffer with zero-rows and same schema as table from + * `create_parquet_with_stats` + */ +template +std::vector create_empty_parquet_with_stats() +{ + auto const non_empty = std::get<0>(create_parquet_with_stats()); + auto const empty = cudf::empty_like(non_empty->view()); + + cudf::io::table_input_metadata output_metadata(empty->view()); + output_metadata.column_metadata[0].set_name("col0"); + output_metadata.column_metadata[1].set_name("col1"); + output_metadata.column_metadata[2].set_name("col2"); + + std::vector buffer; + auto out_opts = + cudf::io::parquet_writer_options::builder(cudf::io::sink_info{&buffer}, empty->view()) + .metadata(std::move(output_metadata)) + .stats_level(cudf::io::statistics_freq::STATISTICS_COLUMN) + .build(); + cudf::io::write_parquet(out_opts); + return buffer; +} + +} // namespace + +struct HybridScanMultifileFiltersTest : public cudf::test::BaseFixture {}; + +TEST_F(HybridScanMultifileFiltersTest, Metadata) +{ + using T = cudf::timestamp_ms; + + // Create two parquet sources, each with 4 row groups and 5000 rows per row + // group + auto constexpr rows_per_row_group = page_size_for_ordered_tests; + auto constexpr num_sources = 2; + + // Build sources with different seeds + std::vector> file_buffers; + file_buffers.reserve(num_sources); + auto constexpr num_concat = 1; + srand(0xbad); + file_buffers.emplace_back(std::get<1>(create_parquet_with_stats())); + srand(0xf00d); + file_buffers.emplace_back(std::get<1>(create_parquet_with_stats())); + + // Filtering AST - col0 < 100 + auto literal_value = + cudf::timestamp_scalar(T(typename T::duration(100)), true, cudf::get_default_stream()); + auto literal = cudf::ast::literal(literal_value); + auto col_ref_0 = cudf::ast::column_name_reference("col0"); + auto filter_expression = cudf::ast::operation(cudf::ast::ast_operator::LESS, col_ref_0, literal); + + // Construct reader from footer bytes + auto inputs = build_multifile_inputs(file_buffers); + + cudf::io::parquet_reader_options options = + cudf::io::parquet_reader_options::builder().filter(filter_expression).build(); + auto const reader = std::make_unique( + cudf::host_span const>{inputs.footer_byte_spans}, options); + + // Get parquet metadata and check + auto parquet_metadata = reader->parquet_metadata(); + ASSERT_EQ(parquet_metadata.size(), num_sources); + for (auto const& meta : parquet_metadata) { + ASSERT_FALSE(meta.row_groups.empty()); + EXPECT_FALSE(meta.row_groups[0].columns[0].offset_index.has_value()); + EXPECT_FALSE(meta.row_groups[0].columns[0].column_index.has_value()); + } + + // Setup page index + auto const page_index_byte_ranges = reader->page_index_byte_range(); + ASSERT_EQ(page_index_byte_ranges.size(), num_sources); + EXPECT_TRUE(std::all_of(page_index_byte_ranges.begin(), + page_index_byte_ranges.end(), + [](auto const& range) { return not range.is_empty(); })); + + std::vector> page_index_buffers; + std::vector> page_index_byte_spans; + page_index_buffers.reserve(num_sources); + page_index_byte_spans.reserve(num_sources); + + auto iter = cuda::zip_iterator(page_index_byte_ranges.begin(), inputs.datasources.begin()); + std::for_each(iter, iter + num_sources, [&](auto const& pair) { + auto const& pgidx_byte_range = cuda::std::get<0>(pair); + auto const& datasource = cuda::std::get<1>(pair); + page_index_buffers.emplace_back( + cudf::io::parquet::fetch_page_index_to_host(*datasource, pgidx_byte_range)); + page_index_byte_spans.emplace_back(*page_index_buffers.back()); + }); + + reader->setup_page_index( + cudf::host_span const>{page_index_byte_spans}); + + // Check if page index is now present in each parquet metadata + parquet_metadata = reader->parquet_metadata(); + for (auto const& meta : parquet_metadata) { + EXPECT_TRUE(meta.row_groups[0].columns[0].offset_index.has_value()); + EXPECT_TRUE(meta.row_groups[0].columns[0].column_index.has_value()); + } + + // Check all row groups + auto input_row_group_indices = reader->all_row_groups(options); + ASSERT_EQ(input_row_group_indices.size(), num_sources); + EXPECT_TRUE(std::all_of( + input_row_group_indices.begin(), input_row_group_indices.end(), [](auto const& rgs) { + return rgs == (std::vector{0, 1, 2, 3}); + })); + + // Set explicit row groups (per-source) via options + options.set_row_groups({{0, 1}, {2, 3}}); + input_row_group_indices = reader->all_row_groups(options); + + // Check if the row groups are set correctly + ASSERT_EQ(input_row_group_indices.size(), num_sources); + EXPECT_EQ(input_row_group_indices[0], (std::vector{0, 1})); + EXPECT_EQ(input_row_group_indices[1], (std::vector{2, 3})); + EXPECT_EQ(reader->total_rows_in_row_groups(input_row_group_indices), + 2 * rows_per_row_group * num_sources); + + // Construct a new reader from a span of existing FileMetaData + auto const reader_with_existing_metadata = + std::make_unique( + cudf::host_span{parquet_metadata}, options); + + // Check if the new metadata is the same as the existing one + auto const new_metadata = reader_with_existing_metadata->parquet_metadata(); + ASSERT_EQ(new_metadata.size(), num_sources); + EXPECT_TRUE(std::all_of(new_metadata.begin(), new_metadata.end(), [&](auto const& meta) { + return meta.row_groups.size() == parquet_metadata.front().row_groups.size(); + })); +} + +TEST_F(HybridScanMultifileFiltersTest, EmptySource) +{ + using T = uint32_t; + + srand(0xc0ffee); + + // Create two parquet source. First one with non-zero rows and the second one with zero rows. + auto constexpr num_sources = 2; + std::vector> file_buffers; + file_buffers.reserve(num_sources); + file_buffers.emplace_back(std::get<1>(create_parquet_with_stats())); + file_buffers.emplace_back(create_empty_parquet_with_stats()); + + auto inputs = build_multifile_inputs(file_buffers); + + cudf::io::parquet_reader_options options = cudf::io::parquet_reader_options::builder().build(); + auto const reader = std::make_unique( + inputs.footer_byte_spans, options); + + // Check parquet metadata + auto const parquet_metadata = reader->parquet_metadata(); + ASSERT_EQ(parquet_metadata.size(), num_sources); + EXPECT_FALSE(parquet_metadata.front().row_groups.empty()); + EXPECT_TRUE(parquet_metadata.back().row_groups.empty()); + + // Check row group indices + auto const all_rgs = reader->all_row_groups(options); + ASSERT_EQ(all_rgs.size(), num_sources); + EXPECT_EQ(all_rgs.front(), (std::vector{0, 1, 2, 3})); + EXPECT_TRUE(all_rgs.back().empty()); + + // Check page index byte ranges + auto const page_index_byte_ranges = reader->page_index_byte_range(); + ASSERT_EQ(page_index_byte_ranges.size(), num_sources); + EXPECT_FALSE(page_index_byte_ranges.front().is_empty()); + EXPECT_TRUE(page_index_byte_ranges.back().is_empty()); +} From f8b90ea171218fba4d46b09ae3332d977950db3f Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 13:53:22 -0700 Subject: [PATCH 05/15] Apply suggestions from code review Co-authored-by: Yunsong Wang <12716979+PointKernel@users.noreply.github.com> --- cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp | 4 ++-- cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp | 4 ++-- .../io/experimental/hybrid_scan_multifile_filters_test.cpp | 2 +- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index 45ea96ab3628..f542173fc23d 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -1,5 +1,5 @@ /* - * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION. * SPDX-License-Identifier: Apache-2.0 */ @@ -19,7 +19,7 @@ #include #include -namespace CUDF_EXPORT cudf { +namespace cudf { namespace io::parquet::experimental::detail { /** * @brief Internal experimental Parquet reader optimized for highly selective filters, called a diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index 9a8879163782..a81aa827aebf 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -114,14 +114,14 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { * * @return Vector of byte ranges of the page index, one per source */ - [[nodiscard]] std::vector page_index_byte_range() const; + [[nodiscard]] std::vector page_index_byte_ranges() const; /** * @brief Get the Parquet file metadata for every source * * @return Vector of file metadata, one per source */ - [[nodiscard]] std::vector parquet_metadata() const; + [[nodiscard]] std::vector parquet_metadatas() const; /** * @brief Setup and populate the page index structs in every source's `FileMetaData` diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp index 0c62cc8b223c..bee1f7edfcc0 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -1,5 +1,5 @@ /* - * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION. * SPDX-License-Identifier: Apache-2.0 */ From 90aef6178439d5225f88e9a9ae702339dd0abeda Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 13:53:53 -0700 Subject: [PATCH 06/15] Update cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp Co-authored-by: Yunsong Wang <12716979+PointKernel@users.noreply.github.com> --- cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index 32819903f530..97592add9ac6 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -1,5 +1,5 @@ /* - * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION. * SPDX-License-Identifier: Apache-2.0 */ From 8876f5711fd067eda4391285ea7c4c09aea7f380 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 21:43:35 +0000 Subject: [PATCH 07/15] Apply suggestions from @PointKernel (thanks!) --- .../cudf/io/experimental/hybrid_scan.hpp | 6 ++---- .../io/experimental/hybrid_scan_multifile.hpp | 18 ++++++++---------- .../io/parquet/experimental/hybrid_scan.cpp | 14 +++++++------- .../experimental/hybrid_scan_helpers.cpp | 18 ++++++++---------- .../experimental/hybrid_scan_helpers.hpp | 2 +- .../parquet/experimental/hybrid_scan_impl.cpp | 12 ++++++------ .../parquet/experimental/hybrid_scan_impl.hpp | 12 ++++++------ .../experimental/hybrid_scan_multifile.cpp | 12 ++++++------ .../hybrid_scan_multifile_filters_test.cpp | 17 ++++++++--------- 9 files changed, 52 insertions(+), 59 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan.hpp b/cpp/include/cudf/io/experimental/hybrid_scan.hpp index 0f7ead8a3dfd..709acfee2804 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan.hpp @@ -19,15 +19,13 @@ #include #include -namespace CUDF_EXPORT cudf { -namespace io::parquet::experimental::detail { +namespace cudf::io::parquet::experimental::detail { /** * @brief Internal experimental Parquet reader optimized for highly selective filters, called a * Hybrid Scan operation. */ class hybrid_scan_reader_impl; -} // namespace io::parquet::experimental::detail -} // namespace CUDF_EXPORT cudf +} // namespace cudf::io::parquet::experimental::detail //! Using `byte_range_info` from cudf::io::text using cudf::io::text::byte_range_info; diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index f542173fc23d..70ff792660bf 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -19,15 +19,13 @@ #include #include -namespace cudf { -namespace io::parquet::experimental::detail { +namespace cudf::io::parquet::experimental::detail { /** * @brief Internal experimental Parquet reader optimized for highly selective filters, called a * Hybrid Scan operation. */ class hybrid_scan_reader_impl; -} // namespace io::parquet::experimental::detail -} // namespace CUDF_EXPORT cudf +} // namespace cudf::io::parquet::experimental::detail //! Using `byte_range_info` from cudf::io::text using cudf::io::text::byte_range_info; @@ -77,25 +75,25 @@ class hybrid_scan_multifile { ~hybrid_scan_multifile(); /** - * @brief Get the per-source Parquet file footer metadata + * @brief Get parquet metadatas for all sources * - * @return Vector of file metadata, one per source + * @return Vector of parquet metadata, one per source */ - [[nodiscard]] std::vector parquet_metadata() const; + [[nodiscard]] std::vector parquet_metadatas() const; /** - * @brief Get the per-source byte range of the page index in each Parquet file + * @brief Get byte ranges of the page index for all sources * * @return Vector of page index byte ranges, one per source */ - [[nodiscard]] std::vector page_index_byte_range() const; + [[nodiscard]] std::vector page_index_byte_ranges() const; /** * @brief Setup the per-source page index within each Parquet file metadata * * @param page_index_bytes Host span of Parquet page index buffer bytes, one per source */ - void setup_page_index( + void setup_page_indexes( cudf::host_span const> page_index_bytes) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 3223ca99c8f8..96ec82d30f18 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -15,28 +15,28 @@ namespace cudf::io::parquet::experimental { hybrid_scan_reader::hybrid_scan_reader(cudf::host_span footer_bytes, parquet_reader_options const& options) + : _impl{std::make_unique( + std::vector>{footer_bytes}, options)} { - auto const footers = std::vector>{footer_bytes}; - _impl = std::make_unique(footers, options); } hybrid_scan_reader::hybrid_scan_reader(FileMetaData const& parquet_metadata, parquet_reader_options const& options) + : _impl{std::make_unique( + std::vector{parquet_metadata}, options)} { - auto const metadatas = std::vector{parquet_metadata}; - _impl = std::make_unique(metadatas, options); } hybrid_scan_reader::~hybrid_scan_reader() = default; [[nodiscard]] text::byte_range_info hybrid_scan_reader::page_index_byte_range() const { - return _impl->page_index_byte_range().front(); + return _impl->page_index_byte_ranges().front(); } [[nodiscard]] FileMetaData hybrid_scan_reader::parquet_metadata() const { - return _impl->parquet_metadata().front(); + return _impl->parquet_metadatas().front(); } void hybrid_scan_reader::setup_page_index(cudf::host_span page_index_bytes) const @@ -44,7 +44,7 @@ void hybrid_scan_reader::setup_page_index(cudf::host_span page_in CUDF_FUNC_RANGE(); auto const per_source = std::vector>{page_index_bytes}; - return _impl->setup_page_index(per_source); + return _impl->setup_page_indexes(per_source); } std::vector hybrid_scan_reader::all_row_groups( diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index e93054073637..4cebe6551da6 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -135,15 +135,15 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, } } -std::vector aggregate_reader_metadata::page_index_byte_range() const +std::vector aggregate_reader_metadata::page_index_byte_ranges() const { std::vector page_index_byte_ranges; page_index_byte_ranges.reserve(per_file_metadata.size()); std::transform(per_file_metadata.begin(), per_file_metadata.end(), std::back_inserter(page_index_byte_ranges), - [](auto const& pfm) -> text::byte_range_info { - auto const& row_groups = pfm.row_groups; + [](auto const& file_metadata) -> text::byte_range_info { + auto const& row_groups = file_metadata.row_groups; if (row_groups.empty() or row_groups.front().columns.empty()) { return {}; } auto const min_offset = row_groups.front().columns.front().column_index_offset; @@ -158,12 +158,12 @@ std::vector aggregate_reader_metadata::page_index_byte_ra return page_index_byte_ranges; } -std::vector aggregate_reader_metadata::parquet_metadata() const +std::vector aggregate_reader_metadata::parquet_metadatas() const { return {per_file_metadata.begin(), per_file_metadata.end()}; } -void aggregate_reader_metadata::setup_page_index( +void aggregate_reader_metadata::setup_page_indexes( cudf::host_span const> page_index_bytes) { CUDF_EXPECTS(page_index_bytes.size() == per_file_metadata.size(), @@ -171,15 +171,13 @@ void aggregate_reader_metadata::setup_page_index( auto iter = cuda::zip_iterator(page_index_bytes.begin(), per_file_metadata.begin()); std::for_each(iter, iter + page_index_bytes.size(), [&](auto const& pair) { - auto const& pgidx_bytes = cuda::std::get<0>(pair); + // Get the page index bytes and file metadata + auto const& [pgidx_bytes, file_metadata] = pair; + auto const& row_groups = file_metadata.row_groups; // Return early if empty page index buffer span if (pgidx_bytes.empty()) { return; } - // Get the file metadata and setup the page index - auto& file_metadata = cuda::std::get<1>(pair); - auto const& row_groups = file_metadata.row_groups; - // Check for empty parquet file CUDF_EXPECTS(not row_groups.empty() and not row_groups.front().columns.empty(), "No column chunks in Parquet schema to read page index for"); diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index a81aa827aebf..2bc10699847e 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -128,7 +128,7 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { * * @param page_index_bytes Host span of Parquet page index buffer bytes, one per source */ - void setup_page_index(cudf::host_span const> page_index_bytes); + void setup_page_indexes(cudf::host_span const> page_index_bytes); /** * @brief Get all available row group indices, one inner vector per source diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 5ef78f19299e..eababea05f4b 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -86,20 +86,20 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl( _extended_metadata = static_cast(_metadata.get()); } -std::vector hybrid_scan_reader_impl::parquet_metadata() const +std::vector hybrid_scan_reader_impl::parquet_metadatas() const { - return _extended_metadata->parquet_metadata(); + return _extended_metadata->parquet_metadatas(); } -std::vector hybrid_scan_reader_impl::page_index_byte_range() const +std::vector hybrid_scan_reader_impl::page_index_byte_ranges() const { - return _extended_metadata->page_index_byte_range(); + return _extended_metadata->page_index_byte_ranges(); } -void hybrid_scan_reader_impl::setup_page_index( +void hybrid_scan_reader_impl::setup_page_indexes( cudf::host_span const> page_index_bytes) const { - _extended_metadata->setup_page_index(page_index_bytes); + _extended_metadata->setup_page_indexes(page_index_bytes); } void hybrid_scan_reader_impl::select_columns(read_columns_mode read_columns_mode, diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 0e85f066d347..21d727f51be3 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -58,19 +58,19 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { parquet_reader_options const& options); /** - * @copydoc cudf::io::experimental::hybrid_scan_multifile::parquet_metadata + * @copydoc cudf::io::experimental::hybrid_scan_multifile::parquet_metadatas */ - [[nodiscard]] std::vector parquet_metadata() const; + [[nodiscard]] std::vector parquet_metadatas() const; /** - * @copydoc cudf::io::experimental::hybrid_scan_multifile::page_index_byte_range + * @copydoc cudf::io::experimental::hybrid_scan_multifile::page_index_byte_ranges */ - [[nodiscard]] std::vector page_index_byte_range() const; + [[nodiscard]] std::vector page_index_byte_ranges() const; /** - * @copydoc cudf::io::experimental::hybrid_scan_multifile::setup_page_index + * @copydoc cudf::io::experimental::hybrid_scan_multifile::setup_page_indexes */ - void setup_page_index( + void setup_page_indexes( cudf::host_span const> page_index_bytes) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index 97592add9ac6..31b3cb5a6443 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -25,21 +25,21 @@ hybrid_scan_multifile::hybrid_scan_multifile(cudf::host_span hybrid_scan_multifile::~hybrid_scan_multifile() = default; -std::vector hybrid_scan_multifile::parquet_metadata() const +std::vector hybrid_scan_multifile::parquet_metadatas() const { - return _impl->parquet_metadata(); + return _impl->parquet_metadatas(); } -std::vector hybrid_scan_multifile::page_index_byte_range() const +std::vector hybrid_scan_multifile::page_index_byte_ranges() const { - return _impl->page_index_byte_range(); + return _impl->page_index_byte_ranges(); } -void hybrid_scan_multifile::setup_page_index( +void hybrid_scan_multifile::setup_page_indexes( cudf::host_span const> page_index_bytes) const { CUDF_FUNC_RANGE(); - _impl->setup_page_index(page_index_bytes); + _impl->setup_page_indexes(page_index_bytes); } std::vector> hybrid_scan_multifile::all_row_groups( diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp index bee1f7edfcc0..1a0c135207f4 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -111,7 +111,7 @@ TEST_F(HybridScanMultifileFiltersTest, Metadata) cudf::host_span const>{inputs.footer_byte_spans}, options); // Get parquet metadata and check - auto parquet_metadata = reader->parquet_metadata(); + auto parquet_metadata = reader->parquet_metadatas(); ASSERT_EQ(parquet_metadata.size(), num_sources); for (auto const& meta : parquet_metadata) { ASSERT_FALSE(meta.row_groups.empty()); @@ -120,7 +120,7 @@ TEST_F(HybridScanMultifileFiltersTest, Metadata) } // Setup page index - auto const page_index_byte_ranges = reader->page_index_byte_range(); + auto const page_index_byte_ranges = reader->page_index_byte_ranges(); ASSERT_EQ(page_index_byte_ranges.size(), num_sources); EXPECT_TRUE(std::all_of(page_index_byte_ranges.begin(), page_index_byte_ranges.end(), @@ -133,18 +133,17 @@ TEST_F(HybridScanMultifileFiltersTest, Metadata) auto iter = cuda::zip_iterator(page_index_byte_ranges.begin(), inputs.datasources.begin()); std::for_each(iter, iter + num_sources, [&](auto const& pair) { - auto const& pgidx_byte_range = cuda::std::get<0>(pair); - auto const& datasource = cuda::std::get<1>(pair); + auto const& [pgidx_byte_range, datasource] = pair; page_index_buffers.emplace_back( cudf::io::parquet::fetch_page_index_to_host(*datasource, pgidx_byte_range)); page_index_byte_spans.emplace_back(*page_index_buffers.back()); }); - reader->setup_page_index( + reader->setup_page_indexes( cudf::host_span const>{page_index_byte_spans}); // Check if page index is now present in each parquet metadata - parquet_metadata = reader->parquet_metadata(); + parquet_metadata = reader->parquet_metadatas(); for (auto const& meta : parquet_metadata) { EXPECT_TRUE(meta.row_groups[0].columns[0].offset_index.has_value()); EXPECT_TRUE(meta.row_groups[0].columns[0].column_index.has_value()); @@ -175,7 +174,7 @@ TEST_F(HybridScanMultifileFiltersTest, Metadata) cudf::host_span{parquet_metadata}, options); // Check if the new metadata is the same as the existing one - auto const new_metadata = reader_with_existing_metadata->parquet_metadata(); + auto const new_metadata = reader_with_existing_metadata->parquet_metadatas(); ASSERT_EQ(new_metadata.size(), num_sources); EXPECT_TRUE(std::all_of(new_metadata.begin(), new_metadata.end(), [&](auto const& meta) { return meta.row_groups.size() == parquet_metadata.front().row_groups.size(); @@ -202,7 +201,7 @@ TEST_F(HybridScanMultifileFiltersTest, EmptySource) inputs.footer_byte_spans, options); // Check parquet metadata - auto const parquet_metadata = reader->parquet_metadata(); + auto const parquet_metadata = reader->parquet_metadatas(); ASSERT_EQ(parquet_metadata.size(), num_sources); EXPECT_FALSE(parquet_metadata.front().row_groups.empty()); EXPECT_TRUE(parquet_metadata.back().row_groups.empty()); @@ -214,7 +213,7 @@ TEST_F(HybridScanMultifileFiltersTest, EmptySource) EXPECT_TRUE(all_rgs.back().empty()); // Check page index byte ranges - auto const page_index_byte_ranges = reader->page_index_byte_range(); + auto const page_index_byte_ranges = reader->page_index_byte_ranges(); ASSERT_EQ(page_index_byte_ranges.size(), num_sources); EXPECT_FALSE(page_index_byte_ranges.front().is_empty()); EXPECT_TRUE(page_index_byte_ranges.back().is_empty()); From a5752762ca11682482f7830bd9e95833a0c934d4 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 21:47:35 +0000 Subject: [PATCH 08/15] Minor --- cpp/src/io/parquet/experimental/hybrid_scan.cpp | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 96ec82d30f18..3bc68cb97752 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -42,9 +42,7 @@ hybrid_scan_reader::~hybrid_scan_reader() = default; void hybrid_scan_reader::setup_page_index(cudf::host_span page_index_bytes) const { CUDF_FUNC_RANGE(); - - auto const per_source = std::vector>{page_index_bytes}; - return _impl->setup_page_indexes(per_source); + return _impl->setup_page_indexes(std::vector>{page_index_bytes}); } std::vector hybrid_scan_reader::all_row_groups( From c7d7cb6da3bf534f9def672c02ce719d51b0631a Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 21:51:32 +0000 Subject: [PATCH 09/15] Minor changes --- .../experimental/hybrid_scan_helpers.cpp | 29 ++++++++++--------- 1 file changed, 16 insertions(+), 13 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 4cebe6551da6..0d29fdfd358c 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -202,7 +202,7 @@ std::vector> aggregate_reader_metadata::all_row_groups( auto const& opts_row_groups = options.get_row_groups(); if (not opts_row_groups.empty()) { CUDF_EXPECTS(opts_row_groups.size() == per_file_metadata.size(), - "Row groups in parquet reader options must have one inner vector per source"); + "Row groups in parquet reader options must specify one vector per data source"); return opts_row_groups; } @@ -224,18 +224,21 @@ size_type aggregate_reader_metadata::total_rows_in_row_groups( { std::size_t total_rows = 0; - std::for_each(cuda::counting_iterator{0}, - cuda::counting_iterator{row_group_indices.size()}, - [&](auto const src_idx) { - auto const& pfm = per_file_metadata[src_idx]; - for (auto const row_group_idx : row_group_indices[src_idx]) { - CUDF_EXPECTS(std::cmp_less(row_group_idx, pfm.row_groups.size()), - "Row group index out of bounds"); - total_rows += pfm.row_groups[row_group_idx].num_rows; - } - }); - CUDF_EXPECTS(std::cmp_less_equal(total_rows, std::numeric_limits::max()), - "Total number of rows exceeds cudf::size_type's limit"); + std::for_each( + cuda::counting_iterator{0}, + cuda::counting_iterator{row_group_indices.size()}, + [&](auto const src_idx) { + auto const& pfm = per_file_metadata[src_idx]; + for (auto const row_group_idx : row_group_indices[src_idx]) { + CUDF_EXPECTS( + std::cmp_greater_equal(row_group_idx, size_type{0}) and + std::cmp_less(row_group_idx, pfm.row_groups.size()), + "Encountered out-of-bounds row group index for data source. Row group index: " + + std::to_string(row_group_idx) + ", Source index: " + std::to_string(src_idx) + + ", Number of row groups: " + std::to_string(pfm.row_groups.size())); + total_rows += pfm.row_groups[row_group_idx].num_rows; + } + }); return static_cast(total_rows); } From 467a628f3fe6176e445ecad6266ac6abb527124d Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 21:53:00 +0000 Subject: [PATCH 10/15] Allow more than 2B rows --- cpp/include/cudf/io/experimental/hybrid_scan.hpp | 2 +- cpp/src/io/parquet/experimental/hybrid_scan.cpp | 2 +- cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp | 4 ++-- cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp | 2 +- cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp | 2 +- cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp | 2 +- 6 files changed, 7 insertions(+), 7 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan.hpp b/cpp/include/cudf/io/experimental/hybrid_scan.hpp index 709acfee2804..980ab9644d3b 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan.hpp @@ -342,7 +342,7 @@ class hybrid_scan_reader { * @param row_group_indices Input row groups indices * @return Total number of top-level rows in the row groups */ - [[nodiscard]] size_type total_rows_in_row_groups( + [[nodiscard]] std::size_t total_rows_in_row_groups( cudf::host_span row_group_indices) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 3bc68cb97752..868813a1b4ed 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -57,7 +57,7 @@ std::vector hybrid_scan_reader::all_row_groups( return _impl->all_row_groups(options).front(); } -size_type hybrid_scan_reader::total_rows_in_row_groups( +std::size_t hybrid_scan_reader::total_rows_in_row_groups( cudf::host_span row_group_indices) const { if (row_group_indices.empty()) { return 0; } diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 0d29fdfd358c..6e12055dacae 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -219,7 +219,7 @@ std::vector> aggregate_reader_metadata::all_row_groups( return row_groups; } -size_type aggregate_reader_metadata::total_rows_in_row_groups( +std::size_t aggregate_reader_metadata::total_rows_in_row_groups( cudf::host_span const> row_group_indices) const { std::size_t total_rows = 0; @@ -240,7 +240,7 @@ size_type aggregate_reader_metadata::total_rows_in_row_groups( } }); - return static_cast(total_rows); + return total_rows; } std::tuple, diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index 2bc10699847e..e65db678c2d1 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -149,7 +149,7 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { * @param row_group_indices Input row groups indices * @return Total number of top-level rows in the row groups */ - [[nodiscard]] size_type total_rows_in_row_groups( + [[nodiscard]] std::size_t total_rows_in_row_groups( cudf::host_span const> row_group_indices) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index eababea05f4b..f2919c64519f 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -206,7 +206,7 @@ std::vector> hybrid_scan_reader_impl::all_row_groups( return _extended_metadata->all_row_groups(options); } -size_type hybrid_scan_reader_impl::total_rows_in_row_groups( +std::size_t hybrid_scan_reader_impl::total_rows_in_row_groups( cudf::host_span const> row_group_indices) const { return _extended_metadata->total_rows_in_row_groups(row_group_indices); diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 21d727f51be3..64eead4463f7 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -82,7 +82,7 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { /** * @copydoc cudf::io::experimental::hybrid_scan::total_rows_in_row_groups */ - [[nodiscard]] size_type total_rows_in_row_groups( + [[nodiscard]] std::size_t total_rows_in_row_groups( cudf::host_span const> row_group_indices) const; /** From 1909146d843f1d376915e03bbe00629710b87f16 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 22:29:44 +0000 Subject: [PATCH 11/15] Minor bug fix --- cpp/src/io/parquet/experimental/page_index_filter.cu | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/cpp/src/io/parquet/experimental/page_index_filter.cu b/cpp/src/io/parquet/experimental/page_index_filter.cu index 21ae60443f3d..7ab305174d3e 100644 --- a/cpp/src/io/parquet/experimental/page_index_filter.cu +++ b/cpp/src/io/parquet/experimental/page_index_filter.cu @@ -994,8 +994,11 @@ thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( "Input row bitmask should be of type BOOL8"); auto const total_rows = total_rows_in_row_groups(row_group_indices); + CUDF_EXPECTS(std::cmp_less_equal(total_rows, std::numeric_limits::max()), + "Total rows in row groups exceed the cudf column size limit", + std::overflow_error); - CUDF_EXPECTS(row_mask_offset + total_rows <= row_mask.size(), + CUDF_EXPECTS(std::cmp_less_equal(row_mask_offset + total_rows, row_mask.size()), "Mismatch in total rows in input row mask and row groups", std::invalid_argument); @@ -1091,8 +1094,8 @@ thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( if constexpr (cuda::std::is_same_v) { if (row_mask.nullable() and row_mask.null_count() > 0) { thrust::for_each(rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()), - cuda::counting_iterator{row_mask_offset}, - cuda::counting_iterator{row_mask_offset + total_rows}, + cuda::counting_iterator(row_mask_offset), + cuda::counting_iterator(row_mask_offset + total_rows), [row_mask = row_mask.template begin(), null_mask = row_mask.null_mask()] __device__(auto const row_idx) { if (not bit_is_set(null_mask, row_idx)) { row_mask[row_idx] = true; } @@ -1126,7 +1129,7 @@ thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( cudf::detail::make_device_uvector_async(host_tree_level_ptrs, stream, mr); // Build Fenwick tree levels (zeroth level is just the row mask itself) - auto prev_level_size = total_rows; + auto prev_level_size = static_cast(total_rows); std::for_each( cuda::counting_iterator{0}, cuda::counting_iterator{num_levels - 1}, From 17fd247d065134b1eeefa3c8358c380ca840c4ba Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Thu, 21 May 2026 22:31:14 +0000 Subject: [PATCH 12/15] Minor --- cpp/src/io/parquet/experimental/page_index_filter.cu | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/cpp/src/io/parquet/experimental/page_index_filter.cu b/cpp/src/io/parquet/experimental/page_index_filter.cu index 7ab305174d3e..0b2aec183b39 100644 --- a/cpp/src/io/parquet/experimental/page_index_filter.cu +++ b/cpp/src/io/parquet/experimental/page_index_filter.cu @@ -1003,7 +1003,8 @@ thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( std::invalid_argument); // Return an empty vector if all rows are invalid or all rows are required - if (row_mask.null_count(row_mask_offset, row_mask_offset + total_rows, stream) == total_rows or + if (std::cmp_equal(row_mask.null_count(row_mask_offset, row_mask_offset + total_rows, stream), + total_rows) or cudf::detail::all_of(row_mask.template begin() + row_mask_offset, row_mask.template begin() + row_mask_offset + total_rows, cuda::std::identity{}, From 1916d997c97062240cfd26dca35236688edcedae Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Fri, 29 May 2026 19:36:08 +0000 Subject: [PATCH 13/15] Minor improvement --- .../experimental/hybrid_scan_helpers.cpp | 35 +++++++++++-------- 1 file changed, 20 insertions(+), 15 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 6e12055dacae..61a7ebd4d56d 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -222,25 +222,30 @@ std::vector> aggregate_reader_metadata::all_row_groups( std::size_t aggregate_reader_metadata::total_rows_in_row_groups( cudf::host_span const> row_group_indices) const { - std::size_t total_rows = 0; + CUDF_EXPECTS(row_group_indices.size() == per_file_metadata.size(), + "Encountered unexpected number of input row group indices", + std::invalid_argument); - std::for_each( + return std::accumulate( cuda::counting_iterator{0}, cuda::counting_iterator{row_group_indices.size()}, - [&](auto const src_idx) { - auto const& pfm = per_file_metadata[src_idx]; - for (auto const row_group_idx : row_group_indices[src_idx]) { - CUDF_EXPECTS( - std::cmp_greater_equal(row_group_idx, size_type{0}) and - std::cmp_less(row_group_idx, pfm.row_groups.size()), - "Encountered out-of-bounds row group index for data source. Row group index: " + - std::to_string(row_group_idx) + ", Source index: " + std::to_string(src_idx) + - ", Number of row groups: " + std::to_string(pfm.row_groups.size())); - total_rows += pfm.row_groups[row_group_idx].num_rows; - } + std::size_t{0}, + [&](auto sum, auto const src_idx) { + auto const& file_metadata = per_file_metadata[src_idx]; + return std::accumulate( + row_group_indices[src_idx].begin(), + row_group_indices[src_idx].end(), + sum, + [&](auto sum, auto const row_group_idx) { + CUDF_EXPECTS( + std::cmp_greater_equal(row_group_idx, size_type{0}) and + std::cmp_less(row_group_idx, file_metadata.row_groups.size()), + "Encountered out-of-bounds row group index for data source. Row group index: " + + std::to_string(row_group_idx) + ", Source index: " + std::to_string(src_idx) + + ", Number of row groups: " + std::to_string(file_metadata.row_groups.size())); + return sum + file_metadata.row_groups[row_group_idx].num_rows; + }); }); - - return total_rows; } std::tuple, From 620615d15196290bebd1dfb6c9b4ce351f341baf Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Fri, 29 May 2026 19:50:08 +0000 Subject: [PATCH 14/15] Minor --- cpp/src/io/parquet/experimental/page_index_filter.cu | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/cpp/src/io/parquet/experimental/page_index_filter.cu b/cpp/src/io/parquet/experimental/page_index_filter.cu index 0b2aec183b39..bb2b89c56c7e 100644 --- a/cpp/src/io/parquet/experimental/page_index_filter.cu +++ b/cpp/src/io/parquet/experimental/page_index_filter.cu @@ -994,13 +994,11 @@ thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( "Input row bitmask should be of type BOOL8"); auto const total_rows = total_rows_in_row_groups(row_group_indices); - CUDF_EXPECTS(std::cmp_less_equal(total_rows, std::numeric_limits::max()), - "Total rows in row groups exceed the cudf column size limit", - std::overflow_error); - CUDF_EXPECTS(std::cmp_less_equal(row_mask_offset + total_rows, row_mask.size()), - "Mismatch in total rows in input row mask and row groups", - std::invalid_argument); + CUDF_EXPECTS( + std::cmp_less_equal(static_cast(row_mask_offset) + total_rows, row_mask.size()), + "Encountered a mismatch in number of rows in the row group pass and the row mask size", + std::invalid_argument); // Return an empty vector if all rows are invalid or all rows are required if (std::cmp_equal(row_mask.null_count(row_mask_offset, row_mask_offset + total_rows, stream), From 3d04fb4d1db6ab870b6245a37724258afa6e83f8 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Fri, 29 May 2026 20:55:26 +0000 Subject: [PATCH 15/15] Address comments --- .../parquet/experimental/hybrid_scan_helpers.cpp | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 61a7ebd4d56d..626ac249b1bd 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -138,7 +138,6 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, std::vector aggregate_reader_metadata::page_index_byte_ranges() const { std::vector page_index_byte_ranges; - page_index_byte_ranges.reserve(per_file_metadata.size()); std::transform(per_file_metadata.begin(), per_file_metadata.end(), std::back_inserter(page_index_byte_ranges), @@ -237,12 +236,13 @@ std::size_t aggregate_reader_metadata::total_rows_in_row_groups( row_group_indices[src_idx].end(), sum, [&](auto sum, auto const row_group_idx) { - CUDF_EXPECTS( - std::cmp_greater_equal(row_group_idx, size_type{0}) and - std::cmp_less(row_group_idx, file_metadata.row_groups.size()), - "Encountered out-of-bounds row group index for data source. Row group index: " + - std::to_string(row_group_idx) + ", Source index: " + std::to_string(src_idx) + - ", Number of row groups: " + std::to_string(file_metadata.row_groups.size())); + CUDF_EXPECTS(std::cmp_greater_equal(row_group_idx, 0) and + std::cmp_less(row_group_idx, file_metadata.row_groups.size()), + std::format("Encountered out-of-bounds row group index for data source. Row " + "group index: {}, Source index: {}, Number of row groups: {}", + row_group_idx, + src_idx, + file_metadata.row_groups.size())); return sum + file_metadata.row_groups[row_group_idx].num_rows; }); });