From 7b8bb435e3559446c4f0b2462a6038cbe2dea269 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Fri, 5 Jun 2026 03:06:49 +0000 Subject: [PATCH 01/10] Add multifile hybrid scan pass construction API --- .../cudf/io/experimental/hybrid_scan.hpp | 4 +- .../io/experimental/hybrid_scan_multifile.hpp | 21 +++ .../io/parquet/experimental/hybrid_scan.cpp | 13 +- .../parquet/experimental/hybrid_scan_impl.cpp | 83 +++++++----- .../parquet/experimental/hybrid_scan_impl.hpp | 20 ++- .../experimental/hybrid_scan_multifile.cpp | 45 +++++++ cpp/src/io/parquet/reader_impl_helpers.hpp | 7 + .../hybrid_scan_multifile_filters_test.cpp | 126 ++++++++++++++++++ 8 files changed, 281 insertions(+), 38 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan.hpp b/cpp/include/cudf/io/experimental/hybrid_scan.hpp index 980ab9644d3b..816850d48c45 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan.hpp @@ -672,8 +672,8 @@ class hybrid_scan_reader { * @throws cudf::logic_error if `row_group_indices` is empty * * @param row_group_indices Input row group indices - * @param pass_read_limit Limit on the amount of memory used for reading and decompressing row - * group data or `0` if there is no limit + * @param pass_read_limit Memory limit to read and decompress row group data, `0` if there is + * no limit (single pass) * * @return Vector of vectors of row group indices, one per constructed pass */ diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index 70ff792660bf..ddd9c04c1095 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -123,6 +123,27 @@ class hybrid_scan_multifile { */ void reset_column_selection() const; + /** + * @brief Partition row groups into passes such that the amount of GPU memory required to read, + * decompress and decode a pass is bounded by the specified limit + * + * Note that the `pass_read_limit` is a hint, not an absolute limit - if a single row group + * cannot fit within the limit given, it will still constitute a pass. The compressed row group + * size is estimated over all columns in each row group (not just the columns selected for + * reading), for conservative estimates. + * + * @throws cudf::logic_error if no row group indices in the input + * + * @param row_group_indices Input row group indices, one per source + * @param pass_read_limit Memory limit to read and decompress row group data, `0` if there is + * no limit (single pass) + * + * @return Vector of per-source row group indices, one per constructed pass + */ + [[nodiscard]] std::vector>> construct_row_group_passes( + cudf::host_span const> row_group_indices, + std::size_t pass_read_limit) const; + private: std::unique_ptr _impl; }; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 868813a1b4ed..de81eaa05caf 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -370,7 +370,18 @@ table_with_metadata hybrid_scan_reader::materialize_all_columns_chunk() const std::vector> hybrid_scan_reader::construct_row_group_passes( cudf::host_span row_group_indices, std::size_t pass_read_limit) const { - return _impl->construct_row_group_passes(row_group_indices, pass_read_limit); + auto const total_row_groups = row_group_indices.size(); + + CUDF_EXPECTS(total_row_groups > 0, "Empty input row group indices encountered"); + + if (pass_read_limit == 0) { return {{row_group_indices.begin(), row_group_indices.end()}}; } + + auto const input_row_group_indices = + std::vector>{{row_group_indices.begin(), row_group_indices.end()}}; + + return _impl + ->construct_row_group_passes(input_row_group_indices, total_row_groups, pass_read_limit) + .first; } bool hybrid_scan_reader::has_next_table_chunk() const { return _impl->has_next_table_chunk(); } diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index f2919c64519f..54caffc41ef8 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -719,40 +719,47 @@ table_with_metadata hybrid_scan_reader_impl::materialize_all_columns_chunk() return result; } -std::vector> hybrid_scan_reader_impl::construct_row_group_passes( - cudf::host_span row_group_indices, std::size_t pass_read_limit) const +std::pair>, std::vector> +hybrid_scan_reader_impl::construct_row_group_passes( + cudf::host_span const> row_group_indices, + std::size_t total_row_groups, + std::size_t pass_read_limit) const { - CUDF_EXPECTS(not row_group_indices.empty(), "Empty input row group indices encountered"); - - // If pass_read_limit is 0 or there is only one row group, return all in a single pass - if (pass_read_limit == 0 or row_group_indices.size() == 1) { - return {{row_group_indices.begin(), row_group_indices.end()}}; - } + CUDF_EXPECTS( + row_group_indices.size() == _extended_metadata->get_num_sources(), + "Encountered a mismatch in the number of row group indices vectors and the number of input " + "datasources", + std::invalid_argument); - // TODO(mh): Need to handle multiple sources in the future - auto constexpr source_index = 0; + CUDF_EXPECTS( + pass_read_limit > 0, "Pass read limit must be greater than 0", std::invalid_argument); - // Construct row group information auto row_groups_info = std::vector{}; - row_groups_info.reserve(row_group_indices.size()); + row_groups_info.reserve(total_row_groups); size_t start_row = 0; - std::transform(row_group_indices.begin(), - row_group_indices.end(), - std::back_inserter(row_groups_info), - [&](auto const& rg_index) { - auto const& row_group = - _extended_metadata->get_row_group(rg_index, source_index); - auto const [compressed_size, total_size, num_rows, max_leaf_values] = - _extended_metadata->get_row_group_properties(row_group); - auto rg_info = row_group_info{.index = rg_index, - .start_row = start_row, - .unadjusted_num_rows = num_rows, - .source_index = source_index, - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}; - start_row += num_rows; - return rg_info; - }); + std::for_each(cuda::counting_iterator(0), + cuda::counting_iterator(row_group_indices.size()), + [&](auto const source_index) { + auto const& src_row_groups = row_group_indices[source_index]; + std::transform( + src_row_groups.begin(), + src_row_groups.end(), + std::back_inserter(row_groups_info), + [&](auto const rg_index) { + auto const& row_group = + _extended_metadata->get_row_group(rg_index, source_index); + auto const [compressed_size, total_size, num_rows, max_leaf_values] = + _extended_metadata->get_row_group_properties(row_group); + auto rg_info = row_group_info{.index = rg_index, + .start_row = start_row, + .unadjusted_num_rows = num_rows, + .source_index = source_index, + .compressed_size = compressed_size, + .max_leaf_values = max_leaf_values}; + start_row += num_rows; + return rg_info; + }); + }); auto const comp_read_limit = static_cast( pass_read_limit * cudf::io::parquet::detail::input_limit_compression_reserve); @@ -764,15 +771,27 @@ std::vector> hybrid_scan_reader_impl::construct_row auto const& offsets = pass_data.pass_row_group_offsets; auto passes = std::vector>{}; passes.reserve(offsets.size() - 1); + auto row_group_source_map = std::vector{}; + auto const has_multiple_sources = row_group_indices.size() > 1; + if (has_multiple_sources) { row_group_source_map.reserve(row_groups_info.size()); } std::transform(offsets.begin(), offsets.end() - 1, offsets.begin() + 1, std::back_inserter(passes), [&](auto const start, auto const end) { - return std::vector{row_group_indices.begin() + start, - row_group_indices.begin() + end}; + auto pass = std::vector{}; + pass.reserve(end - start); + std::for_each(row_groups_info.begin() + start, + row_groups_info.begin() + end, + [&](auto const& rg_info) { + pass.emplace_back(rg_info.index); + if (has_multiple_sources) { + row_group_source_map.emplace_back(rg_info.source_index); + } + }); + return pass; }); - return passes; + return {std::move(passes), std::move(row_group_source_map)}; } bool hybrid_scan_reader_impl::has_next_table_chunk() diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 64eead4463f7..1d565c64ef22 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -25,6 +25,7 @@ #include #include +#include #include namespace cudf::io::parquet::experimental::detail { @@ -275,10 +276,23 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { [[nodiscard]] table_with_metadata materialize_all_columns_chunk(); /** - * @copydoc cudf::io::experimental::hybrid_scan_reader::construct_row_group_passes + * @brief Partition per-source row groups into read passes + * + * @param row_group_indices Input row group indices, one per source + * @param total_row_groups Total number of row groups across all sources + * @param pass_read_limit Memory limit to read and decompress row + * group data + * + * @return Pair of a vector of flattened row group passes and a source index map. The source index + * map is empty for single source input + * + * @throws std::invalid_argument if @p pass_read_limit` is `0`or if @p row_group_indices.size() is + * not equal to the number of input datasources */ - [[nodiscard]] std::vector> construct_row_group_passes( - cudf::host_span row_group_indices, std::size_t pass_read_limit) const; + [[nodiscard]] std::pair>, std::vector> + construct_row_group_passes(cudf::host_span const> row_group_indices, + std::size_t total_row_groups, + std::size_t pass_read_limit) const; /** * @copydoc cudf::io::experimental::hybrid_scan::has_next_table_chunk diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index 31b3cb5a6443..c10ec5f018c3 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -7,6 +7,9 @@ #include #include +#include + +#include namespace cudf::io::parquet::experimental { @@ -57,4 +60,46 @@ size_type hybrid_scan_multifile::total_rows_in_row_groups( void hybrid_scan_multifile::reset_column_selection() const { _impl->reset_column_selection(); } +std::vector>> hybrid_scan_multifile::construct_row_group_passes( + cudf::host_span const> row_group_indices, + std::size_t pass_read_limit) const +{ + auto const total_row_groups = + std::accumulate(row_group_indices.begin(), + row_group_indices.end(), + std::size_t{0}, + [](auto sum, auto const& rgs) { return sum + rgs.size(); }); + CUDF_EXPECTS(total_row_groups > 0, "Empty input row group indices encountered"); + + if (pass_read_limit == 0) { + return { + std::vector>{row_group_indices.begin(), row_group_indices.end()}}; + } + + auto [passes, source_map] = + _impl->construct_row_group_passes(row_group_indices, total_row_groups, pass_read_limit); + + auto source_passes = std::vector>>{}; + source_passes.reserve(passes.size()); + + if (row_group_indices.size() == 1) { + for (auto& pass : passes) { + source_passes.emplace_back(); + source_passes.back().push_back(std::move(pass)); + } + return source_passes; + } + + auto source_map_it = source_map.begin(); + for (auto const& pass : passes) { + auto source_pass = std::vector>(row_group_indices.size()); + for (auto const row_group_index : pass) { + source_pass[*source_map_it++].push_back(row_group_index); + } + source_passes.push_back(std::move(source_pass)); + } + + return source_passes; +} + } // namespace cudf::io::parquet::experimental diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index a3238d66adbe..c9c42a1bd52c 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -411,6 +411,13 @@ class aggregate_reader_metadata { */ [[nodiscard]] auto get_num_row_groups() const { return num_row_groups; } + /** + * @brief Get total number of sources + * + * @return Total number of sources + */ + [[nodiscard]] auto get_num_sources() const { return per_file_metadata.size(); } + /** * @brief Get the number of row groups per file * 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 1a0c135207f4..09f89f103b09 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -8,6 +8,7 @@ #include #include +#include #include #include #include @@ -16,8 +17,12 @@ #include #include +#include #include #include +#include +#include +#include #include namespace { @@ -217,4 +222,125 @@ TEST_F(HybridScanMultifileFiltersTest, EmptySource) 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()); + + auto const passes = reader->construct_row_group_passes(all_rgs, 1); + ASSERT_EQ(passes.size(), all_rgs.front().size()); + for (auto const& pass : passes) { + ASSERT_EQ(pass.size(), num_sources); + ASSERT_EQ(pass.front().size(), 1); + EXPECT_TRUE(pass.back().empty()); + } +} + +TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) +{ + using T = uint32_t; + + srand(0xced); + + 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(std::get<1>(create_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); + + auto const all_rgs = reader->all_row_groups(options); + ASSERT_EQ(all_rgs.size(), num_sources); + EXPECT_TRUE(std::all_of(all_rgs.begin(), all_rgs.end(), [](auto const& rgs) { + return rgs == (std::vector{0, 1, 2, 3}); + })); + + { + auto invalid_rgs = all_rgs; + invalid_rgs.pop_back(); + EXPECT_THROW(static_cast(reader->construct_row_group_passes(invalid_rgs, 0)), + std::invalid_argument); + } + + { + auto const passes = reader->construct_row_group_passes(all_rgs, 0); + ASSERT_EQ(passes.size(), 1); + ASSERT_EQ(passes.front().size(), num_sources); + EXPECT_EQ(passes.front(), all_rgs); + } + + { + auto const passes = reader->construct_row_group_passes(all_rgs, 10'000); + ASSERT_GT(passes.size(), 1); + auto const pass_num_row_groups = [](auto const& pass) { + return std::accumulate( + pass.begin(), pass.end(), std::size_t{0}, [](auto sum, auto const& rgs) { + return sum + rgs.size(); + }); + }; + EXPECT_TRUE(std::any_of(passes.begin(), passes.end(), [&](auto const& pass) { + return pass_num_row_groups(pass) > 1; + })); + + auto flattened = std::vector>{}; + for (auto const& pass : passes) { + ASSERT_EQ(pass.size(), num_sources); + for (auto source_index = std::size_t{0}; source_index < pass.size(); ++source_index) { + std::transform(pass[source_index].begin(), + pass[source_index].end(), + std::back_inserter(flattened), + [source_index](auto const rg_index) { + return std::pair{source_index, rg_index}; + }); + } + } + + auto expected = std::vector>{}; + for (auto source_index = std::size_t{0}; source_index < all_rgs.size(); ++source_index) { + std::transform(all_rgs[source_index].begin(), + all_rgs[source_index].end(), + std::back_inserter(expected), + [source_index](auto const rg_index) { + return std::pair{source_index, rg_index}; + }); + } + EXPECT_EQ(flattened, expected); + } +} + +TEST_F(HybridScanMultifileFiltersTest, RowGroupPassesSingleSourceParity) +{ + using T = uint32_t; + + srand(0xced); + + std::vector> file_buffers; + file_buffers.emplace_back(std::get<1>(create_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 multifile_reader = + std::make_unique( + inputs.footer_byte_spans, options); + auto const single_file_reader = + std::make_unique( + inputs.footer_byte_spans.front(), options); + + auto const all_rgs = multifile_reader->all_row_groups(options); + ASSERT_EQ(all_rgs.size(), 1); + auto constexpr pass_read_limit = std::size_t{10'000}; + auto const multifile_passes = + multifile_reader->construct_row_group_passes(all_rgs, pass_read_limit); + auto const single_file_passes = + single_file_reader->construct_row_group_passes(all_rgs.front(), pass_read_limit); + + auto projected_passes = std::vector>{}; + projected_passes.reserve(multifile_passes.size()); + for (auto const& pass : multifile_passes) { + ASSERT_EQ(pass.size(), 1); + projected_passes.push_back(pass.front()); + } + EXPECT_EQ(projected_passes, single_file_passes); } From 4b45eefa8ccd4d55ba20087db99c7a208435648f Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Mon, 8 Jun 2026 23:33:11 +0000 Subject: [PATCH 02/10] Address review comments --- .../io/parquet/experimental/hybrid_scan.cpp | 5 ++-- .../parquet/experimental/hybrid_scan_impl.cpp | 6 +++++ .../parquet/experimental/hybrid_scan_impl.hpp | 4 +-- .../experimental/hybrid_scan_multifile.cpp | 10 +++---- .../hybrid_scan_multifile_filters_test.cpp | 26 +++++++++++++++++++ 5 files changed, 40 insertions(+), 11 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index de81eaa05caf..b21c6c715ada 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -372,9 +372,8 @@ std::vector> hybrid_scan_reader::construct_row_grou { auto const total_row_groups = row_group_indices.size(); - CUDF_EXPECTS(total_row_groups > 0, "Empty input row group indices encountered"); - - if (pass_read_limit == 0) { return {{row_group_indices.begin(), row_group_indices.end()}}; } + CUDF_EXPECTS( + total_row_groups > 0, "Empty input row group indices encountered", std::invalid_argument); auto const input_row_group_indices = std::vector>{{row_group_indices.begin(), row_group_indices.end()}}; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 54caffc41ef8..4f63f2998845 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -731,6 +731,12 @@ hybrid_scan_reader_impl::construct_row_group_passes( "datasources", std::invalid_argument); + if (pass_read_limit == 0) { + return { + std::vector>{row_group_indices.begin(), row_group_indices.end()}, + std::vector{}}; + } + CUDF_EXPECTS( pass_read_limit > 0, "Pass read limit must be greater than 0", std::invalid_argument); diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 1d565c64ef22..78136f367c95 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -286,8 +286,8 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { * @return Pair of a vector of flattened row group passes and a source index map. The source index * map is empty for single source input * - * @throws std::invalid_argument if @p pass_read_limit` is `0`or if @p row_group_indices.size() is - * not equal to the number of input datasources + * @throws std::invalid_argument if @p row_group_indices.size() is empty or not equal to the + * number of input datasources */ [[nodiscard]] std::pair>, std::vector> construct_row_group_passes(cudf::host_span const> row_group_indices, diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index c10ec5f018c3..e660a7ed9b73 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -69,16 +69,14 @@ std::vector>> hybrid_scan_multifile::construc row_group_indices.end(), std::size_t{0}, [](auto sum, auto const& rgs) { return sum + rgs.size(); }); - CUDF_EXPECTS(total_row_groups > 0, "Empty input row group indices encountered"); - - if (pass_read_limit == 0) { - return { - std::vector>{row_group_indices.begin(), row_group_indices.end()}}; - } + CUDF_EXPECTS( + total_row_groups > 0, "Empty input row group indices encountered", std::invalid_argument); auto [passes, source_map] = _impl->construct_row_group_passes(row_group_indices, total_row_groups, pass_read_limit); + if (pass_read_limit == 0) { return {passes}; } + auto source_passes = std::vector>>{}; source_passes.reserve(passes.size()); 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 09f89f103b09..e5f405497f82 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -256,6 +256,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) return rgs == (std::vector{0, 1, 2, 3}); })); + // Invalid row group indices => throw error { auto invalid_rgs = all_rgs; invalid_rgs.pop_back(); @@ -263,6 +264,16 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) std::invalid_argument); } + // Empty row group indices => throw error + { + auto const empty_rgs = std::vector>(num_sources); + EXPECT_THROW(static_cast(reader->construct_row_group_passes(empty_rgs, 0)), + std::invalid_argument); + EXPECT_THROW(static_cast(reader->construct_row_group_passes(empty_rgs, 1)), + std::invalid_argument); + } + + // Zero pass read limit => single pass with all row groups { auto const passes = reader->construct_row_group_passes(all_rgs, 0); ASSERT_EQ(passes.size(), 1); @@ -270,6 +281,21 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) EXPECT_EQ(passes.front(), all_rgs); } + // Small pass read limit => each row group in its own pass + { + auto const passes = reader->construct_row_group_passes(all_rgs, 1); + ASSERT_EQ(passes.size(), num_sources * all_rgs.front().size()); + for (auto const& pass : passes) { + ASSERT_EQ(pass.size(), num_sources); + auto const pass_num_row_groups = + std::accumulate(pass.begin(), pass.end(), std::size_t{0}, [](auto sum, auto const& rgs) { + return sum + rgs.size(); + }); + EXPECT_EQ(pass_num_row_groups, 1); + } + } + + // Large pass read limit => multiple passes { auto const passes = reader->construct_row_group_passes(all_rgs, 10'000); ASSERT_GT(passes.size(), 1); From 5c1bc3d99076f877d972f270c85c899b015f4e95 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 9 Jun 2026 00:01:39 +0000 Subject: [PATCH 03/10] style --- .../hybrid_scan_multifile_filters_test.cpp | 24 +++++++++---------- 1 file changed, 11 insertions(+), 13 deletions(-) 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 98ebdb6e4b19..e358a8e9853f 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -314,23 +314,21 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) for (auto const& pass : passes) { ASSERT_EQ(pass.size(), num_sources); for (auto source_index = std::size_t{0}; source_index < pass.size(); ++source_index) { - std::transform(pass[source_index].begin(), - pass[source_index].end(), - std::back_inserter(flattened), - [source_index](auto const rg_index) { - return std::pair{source_index, rg_index}; - }); + std::transform( + pass[source_index].begin(), + pass[source_index].end(), + std::back_inserter(flattened), + [source_index](auto const rg_index) { return std::pair{source_index, rg_index}; }); } } auto expected = std::vector>{}; for (auto source_index = std::size_t{0}; source_index < all_rgs.size(); ++source_index) { - std::transform(all_rgs[source_index].begin(), - all_rgs[source_index].end(), - std::back_inserter(expected), - [source_index](auto const rg_index) { - return std::pair{source_index, rg_index}; - }); + std::transform( + all_rgs[source_index].begin(), + all_rgs[source_index].end(), + std::back_inserter(expected), + [source_index](auto const rg_index) { return std::pair{source_index, rg_index}; }); } EXPECT_EQ(flattened, expected); } @@ -355,7 +353,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPassesSingleSourceParity) std::make_unique( inputs.footer_byte_spans.front(), options); - auto const all_rgs = multifile_reader->all_row_groups(options); + auto const all_rgs = multifile_reader->all_row_groups(options); ASSERT_EQ(all_rgs.size(), 1); auto constexpr pass_read_limit = std::size_t{10'000}; auto const multifile_passes = From f28611570e6c179f76bc15fdaa4b8c36edd67c92 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 9 Jun 2026 01:58:12 +0000 Subject: [PATCH 04/10] Update docstrings --- cpp/include/cudf/io/experimental/hybrid_scan.hpp | 2 +- .../cudf/io/experimental/hybrid_scan_multifile.hpp | 2 +- cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp | 10 ++++++---- cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp | 6 +++--- 4 files changed, 11 insertions(+), 9 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan.hpp b/cpp/include/cudf/io/experimental/hybrid_scan.hpp index 816850d48c45..bab4bf3fb93e 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan.hpp @@ -669,7 +669,7 @@ class hybrid_scan_reader { * size is estimated over all columns in each row group (not just the columns selected for * reading), for conservative estimates. * - * @throws cudf::logic_error if `row_group_indices` is empty + * @throws std::invalid_argument if no row group indices in the input * * @param row_group_indices Input row group indices * @param pass_read_limit Memory limit to read and decompress row group data, `0` if there is diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index 8d99912bd250..4d6b0750d725 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -176,7 +176,7 @@ class hybrid_scan_multifile { * size is estimated over all columns in each row group (not just the columns selected for * reading), for conservative estimates. * - * @throws cudf::logic_error if no row group indices in the input + * @throws std::invalid_argument if no row group indices in the input * * @param row_group_indices Input row group indices, one per source * @param pass_read_limit Memory limit to read and decompress row group data, `0` if there is diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 6330c2b1e149..665d431593fa 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -730,10 +730,12 @@ hybrid_scan_reader_impl::construct_row_group_passes( std::size_t pass_read_limit) const { CUDF_EXPECTS( - row_group_indices.size() == _extended_metadata->get_num_sources(), - "Encountered a mismatch in the number of row group indices vectors and the number of input " - "datasources", - std::invalid_argument); + total_row_groups > 0, "Empty input row group indices encountered", std::invalid_argument); + + CUDF_EXPECTS(row_group_indices.size() == _extended_metadata->get_num_sources(), + "Mismatch in the number of row group indices vectors and the number of input " + "datasources", + std::invalid_argument); if (pass_read_limit == 0) { return { diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 0489f46670ea..188360bbafda 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -278,6 +278,9 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { /** * @brief Partition per-source row groups into read passes * + * @throws std::invalid_argument if @p row_group_indices.size() is all empty or not equal to the + * number of input datasources + * * @param row_group_indices Input row group indices, one per source * @param total_row_groups Total number of row groups across all sources * @param pass_read_limit Memory limit to read and decompress row @@ -285,9 +288,6 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { * * @return Pair of a vector of flattened row group passes and a source index map. The source index * map is empty for single source input - * - * @throws std::invalid_argument if @p row_group_indices.size() is empty or not equal to the - * number of input datasources */ [[nodiscard]] std::pair>, std::vector> construct_row_group_passes(cudf::host_span const> row_group_indices, From 987017be37b99ae3fdf5c5660789f85bc110dc7c Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 10 Jun 2026 00:03:59 +0000 Subject: [PATCH 05/10] Minor --- .../hybrid_scan_multifile_filters_test.cpp | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) 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 e358a8e9853f..6aad15b26995 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -237,14 +237,16 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) { using T = uint32_t; - srand(0xced); - 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(std::get<1>(create_parquet_with_stats())); - + std::transform(cuda::counting_iterator(0), + cuda::counting_iterator(num_sources), + std::back_inserter(file_buffers), + [&](auto i) { + srand(0xced + i); + return std::get<1>(create_parquet_with_stats()); + }); auto inputs = build_multifile_inputs(file_buffers); cudf::io::parquet_reader_options options = cudf::io::parquet_reader_options::builder().build(); @@ -336,7 +338,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) TEST_F(HybridScanMultifileFiltersTest, RowGroupPassesSingleSourceParity) { - using T = uint32_t; + using T = cudf::duration_ms; srand(0xced); From 7dff464372c67de67e6c5d651d748244dcab161f Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Sat, 13 Jun 2026 01:14:55 +0000 Subject: [PATCH 06/10] Style fix --- cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp | 2 +- .../io/experimental/hybrid_scan_multifile_filters_test.cpp | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index c0b87748f20a..52b4ce96b862 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -227,7 +227,7 @@ class hybrid_scan_multifile { parquet_reader_options const& options, rmm::cuda_stream_view stream, rmm::device_async_resource_ref mr) const; - + /* * @brief Partition row groups into passes such that the amount of GPU memory required to read, * decompress and decode a pass is bounded by the specified limit 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 6d02895286c1..491dad476941 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -10,8 +10,8 @@ #include #include -#include #include +#include #include #include #include From 56bccf9e4ff21a095de6bda08d8e8bb85d6fa683 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Mon, 15 Jun 2026 16:26:15 +0000 Subject: [PATCH 07/10] style --- cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp | 2 +- 1 file changed, 1 insertion(+), 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 52b4ce96b862..635b4cef5fea 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -228,7 +228,7 @@ class hybrid_scan_multifile { rmm::cuda_stream_view stream, rmm::device_async_resource_ref mr) const; - /* + /** * @brief Partition row groups into passes such that the amount of GPU memory required to read, * decompress and decode a pass is bounded by the specified limit * From 9b2a107fb7bc9c67a9eb2aba884832cde478717c Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Mon, 15 Jun 2026 16:50:15 +0000 Subject: [PATCH 08/10] Minor merge conflict --- .../io/experimental/hybrid_scan_multifile_filters_test.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 491dad476941..29e547c969d1 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -249,7 +249,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) srand(0xced + i); return std::get<1>(create_parquet_with_stats()); }); - auto inputs = build_multifile_inputs(file_buffers); + auto inputs = multifile_inputs(build_source_info(file_buffers)); cudf::io::parquet_reader_options options = cudf::io::parquet_reader_options::builder().build(); auto const reader = std::make_unique( @@ -347,7 +347,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPassesSingleSourceParity) std::vector> file_buffers; file_buffers.emplace_back(std::get<1>(create_parquet_with_stats())); - auto inputs = build_multifile_inputs(file_buffers); + auto inputs = multifile_inputs(build_source_info(file_buffers)); cudf::io::parquet_reader_options options = cudf::io::parquet_reader_options::builder().build(); auto const multifile_reader = From bba5b5554670d9a9e554b682eb23999274d1cd14 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Mon, 15 Jun 2026 17:41:51 +0000 Subject: [PATCH 09/10] Minor --- cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp | 1 - 1 file changed, 1 deletion(-) diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp index 1273085f32a1..b3c927f320a0 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp @@ -120,7 +120,6 @@ void test_hybrid_scan_multifile(std::vector const& columns, stream, mr); tasks.get(); - (void)column_chunk_buffers; auto column_chunk_data = std::vector>{}; for (auto const& source_column_chunks : column_chunks_per_source) { From 285234e4ab8faf969f348e8a6dfcf7eb7c42e41f Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Tue, 16 Jun 2026 00:58:00 +0000 Subject: [PATCH 10/10] Fix the exception type --- python/pylibcudf/tests/io/test_experimental_hybrid_scan.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py index 51d332e6c2f5..7ee3e540c6ec 100644 --- a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py +++ b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py @@ -679,7 +679,9 @@ def test_hybrid_scan_construct_row_group_passes( assert all(passes) # Empty input row groups raise an error - with pytest.raises(RuntimeError): + with pytest.raises( + ValueError, match="Empty input row group indices encountered" + ): simple_hybrid_scan_reader.construct_row_group_passes( [], pass_read_limit )