From 26ec29f35fde79c7d66d01ca1312e1b901edfde0 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Fri, 24 Jul 2026 18:08:22 -0700 Subject: [PATCH 1/9] Size Parquet reader passes by the selected columns only Pass construction sized each row group by summing `total_compressed_size` over *every* column chunk in the row group, via `aggregate_reader_metadata::get_row_group_properties()`. When only a few columns of a wide schema are read, that over-estimates the pass footprint by roughly the projection ratio, so the reader splits into far more passes than the `pass_read_limit` requires. Each extra pass costs another `read_compressed_data()` round trip, another page-header decode, another dictionary decompression, and smaller decode batches. This was also internally inconsistent: the pass *partition* used the all-columns size while `pass.base_mem_size` - which derives the *subpass* budget - already used the selected columns only. Decouple `compute_row_group_passes()` from `row_group_info` by introducing `row_group_pass_size_info`, and add `make_row_group_pass_size_info()`, which reduces over the `ColumnChunkDesc` array that `create_global_chunk_info()` has already built. Those descriptors exist only for the selected leaf columns, so the resulting per-row-group size is exact rather than estimated, and needs no new metadata plumbing. `max_leaf_values` narrows to the selected columns for the same reason: it guards the `size_type` limit on decoded column length, and unselected columns are never dremel-decoded and never allocate an output buffer. `row_group_info::compressed_size` / `max_leaf_values` had no remaining readers, so drop them along with the now-dead `total_size` element of `get_row_group_properties()`'s return. The hybrid scan reader keeps its existing all-columns estimate here; making its planning API column-aware is a follow-up. Co-Authored-By: Claude Opus 5 --- .../parquet/experimental/hybrid_scan_impl.cpp | 53 ++++++------- cpp/src/io/parquet/reader_impl_chunking.cu | 9 ++- .../io/parquet/reader_impl_chunking_utils.cu | 53 ++++++++++--- .../io/parquet/reader_impl_chunking_utils.cuh | 29 ++++++- cpp/src/io/parquet/reader_impl_helpers.cpp | 16 +--- cpp/src/io/parquet/reader_impl_helpers.hpp | 27 +++++-- cpp/tests/io/parquet_chunked_reader_test.cu | 79 +++++++++++++++++++ 7 files changed, 202 insertions(+), 64 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index ffacc6796301..30a144db1bc4 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -800,38 +800,35 @@ hybrid_scan_reader_impl::construct_row_group_passes( CUDF_EXPECTS( pass_read_limit > 0, "Pass read limit must be greater than 0", std::invalid_argument); - auto row_groups_info = std::vector{}; - row_groups_info.reserve(total_row_groups); + // Row group (index, source index) pairs, flattened across sources, parallel to `row_group_sizes` + auto row_group_ids = std::vector>{}; + auto row_group_sizes = std::vector{}; + row_group_ids.reserve(total_row_groups); + row_group_sizes.reserve(total_row_groups); + size_t start_row = 0; 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; - }); + for (auto const rg_index : row_group_indices[source_index]) { + auto const& row_group = + _extended_metadata->get_row_group(rg_index, source_index); + auto const [compressed_size, num_rows, max_leaf_values] = + _extended_metadata->get_row_group_properties(row_group); + row_group_ids.emplace_back(rg_index, source_index); + row_group_sizes.push_back({.start_row = start_row, + .num_rows = num_rows, + .compressed_size = compressed_size, + .max_leaf_values = max_leaf_values}); + start_row += num_rows; + } }); auto const comp_read_limit = static_cast( pass_read_limit * cudf::io::parquet::detail::input_limit_compression_reserve); auto const pass_data = - cudf::io::parquet::detail::compute_row_group_passes(row_groups_info, comp_read_limit, 0); + cudf::io::parquet::detail::compute_row_group_passes(row_group_sizes, comp_read_limit, 0); // Convert offset-based pass boundaries back to vectors of row group indices auto const& offsets = pass_data.pass_row_group_offsets; @@ -839,7 +836,7 @@ hybrid_scan_reader_impl::construct_row_group_passes( 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()); } + if (has_multiple_sources) { row_group_source_map.reserve(row_group_ids.size()); } std::transform(offsets.begin(), offsets.end() - 1, offsets.begin() + 1, @@ -847,12 +844,12 @@ hybrid_scan_reader_impl::construct_row_group_passes( [&](auto const start, auto const 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); + std::for_each(row_group_ids.begin() + start, + row_group_ids.begin() + end, + [&](auto const& [rg_index, source_index]) { + pass.emplace_back(rg_index); if (has_multiple_sources) { - row_group_source_map.emplace_back(rg_info.source_index); + row_group_source_map.emplace_back(source_index); } }); return pass; diff --git a/cpp/src/io/parquet/reader_impl_chunking.cu b/cpp/src/io/parquet/reader_impl_chunking.cu index c2a53826bb05..8acc978c3173 100644 --- a/cpp/src/io/parquet/reader_impl_chunking.cu +++ b/cpp/src/io/parquet/reader_impl_chunking.cu @@ -543,8 +543,15 @@ void reader_impl::compute_input_passes(read_mode mode) ? static_cast(_input_pass_read_limit * input_limit_compression_reserve) : std::numeric_limits::max(); + // Size each row group by the column chunks we are actually going to read. + // `create_global_chunk_info()` has already built one `ColumnChunkDesc` per (row group, selected + // input column), so this is an exact measure of the compressed bytes a pass will hold rather than + // an estimate over the whole row group. + auto const row_group_sizes = + make_row_group_pass_size_info(row_groups_info, _file_itm_data.chunks, _input_columns.size()); + auto pass_data = - compute_row_group_passes(row_groups_info, comp_read_limit, _file_itm_data.global_skip_rows); + compute_row_group_passes(row_group_sizes, comp_read_limit, _file_itm_data.global_skip_rows); _file_itm_data.input_pass_row_group_offsets = std::move(pass_data.pass_row_group_offsets); _file_itm_data.input_pass_start_row_count = std::move(pass_data.pass_start_row_counts); diff --git a/cpp/src/io/parquet/reader_impl_chunking_utils.cu b/cpp/src/io/parquet/reader_impl_chunking_utils.cu index ebca8918171c..2430a5be8f23 100644 --- a/cpp/src/io/parquet/reader_impl_chunking_utils.cu +++ b/cpp/src/io/parquet/reader_impl_chunking_utils.cu @@ -953,9 +953,40 @@ rmm::device_uvector compute_level_decode_sizes(device_span row_groups_info, - std::size_t comp_read_limit, - int64_t skip_rows) +std::vector make_row_group_pass_size_info( + cudf::host_span row_groups_info, + cudf::host_span chunks, + size_t num_input_columns) +{ + CUDF_EXPECTS(chunks.size() == row_groups_info.size() * num_input_columns, + "Mismatch between the number of column chunk descriptors and row groups"); + + std::vector row_group_sizes; + row_group_sizes.reserve(row_groups_info.size()); + + for (size_t rg_idx = 0; rg_idx < row_groups_info.size(); rg_idx++) { + auto const& rgi = row_groups_info[rg_idx]; + auto const chunks_begin = chunks.begin() + (rg_idx * num_input_columns); + auto const chunks_end = chunks_begin + num_input_columns; + size_t compressed_size = 0; + size_t max_leaf_values = 0; + std::for_each(chunks_begin, chunks_end, [&](auto const& chunk) { + compressed_size += chunk.compressed_size; + max_leaf_values = std::max(max_leaf_values, chunk.num_values); + }); + row_group_sizes.push_back({.start_row = rgi.start_row, + .num_rows = rgi.unadjusted_num_rows, + .compressed_size = compressed_size, + .max_leaf_values = max_leaf_values}); + } + + return row_group_sizes; +} + +row_group_pass_data compute_row_group_passes( + cudf::host_span row_group_sizes, + std::size_t comp_read_limit, + int64_t skip_rows) { auto constexpr max_rows_per_pass = static_cast(std::numeric_limits::max()); @@ -970,14 +1001,14 @@ row_group_pass_data compute_row_group_passes(cudf::host_span 1) - ? (rgi.start_row + rgi.unadjusted_num_rows - skip_rows) - : rgi.unadjusted_num_rows; + auto const row_group_rows = (skip_rows and row_group_sizes.size() > 1) + ? (rgi.start_row + rgi.num_rows - skip_rows) + : rgi.num_rows; auto const compressed_rg_size = rgi.compressed_size; auto const row_group_leaf_values = rgi.max_leaf_values; @@ -999,7 +1030,7 @@ row_group_pass_data compute_row_group_passes(cudf::host_span row_groups_info, - std::size_t comp_read_limit, - int64_t skip_rows); +row_group_pass_data compute_row_group_passes( + cudf::host_span row_group_sizes, + std::size_t comp_read_limit, + int64_t skip_rows); + +/** + * @brief Computes per-row-group pass sizing properties from already-built column chunk descriptors + * + * @p chunks is expected to be laid out as `[row_group][input_column]`, i.e. the layout produced by + * `reader_impl::create_global_chunk_info()`. Because those descriptors only exist for the selected + * columns, the resulting sizes are exact rather than estimated. + * + * @param row_groups_info Span of row group metadata + * @param chunks Column chunk descriptors for the selected columns of every row group + * @param num_input_columns Number of selected (leaf) input columns + * @return Per-row-group pass sizing properties, one entry per row group + */ +std::vector make_row_group_pass_size_info( + cudf::host_span row_groups_info, + cudf::host_span chunks, + size_t num_input_columns); } // namespace cudf::io::parquet::detail diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index ed60a4db0e2d..4802deabc653 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -1384,7 +1384,7 @@ std::vector aggregate_reader_metadata::get_pandas_index_names() con return names; } -std::tuple aggregate_reader_metadata::get_row_group_properties( +std::tuple aggregate_reader_metadata::get_row_group_properties( RowGroup const& row_group) const { auto const compressed_size = std::transform_reduce( @@ -1394,8 +1394,6 @@ std::tuple aggregate_reader_metadata::get_row_gr std::plus<>(), [](auto const& colchunk) { return colchunk.meta_data.total_compressed_size; }); - auto const total_size = compressed_size + row_group.total_byte_size; - size_t const max_leaf_values = row_group.columns.empty() ? 0 @@ -1406,7 +1404,7 @@ std::tuple aggregate_reader_metadata::get_row_gr }) ->meta_data.num_values; - return {compressed_size, total_size, static_cast(row_group.num_rows), max_leaf_values}; + return {compressed_size, static_cast(row_group.num_rows), max_leaf_values}; } std::tuple(src_idx), - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}); + .unadjusted_num_rows = static_cast(rg.num_rows), + .source_index = static_cast(src_idx)}); // If page-level indexes are present, then collect extra chunk and page // info. The page indexes rely on absolute row numbers - not adjusted for diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index 40c70babf2b8..2d83d4e95770 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -70,8 +70,6 @@ struct row_group_info { size_t source_start_row; // file-local start row of this row group within its source file size_t unadjusted_num_rows; // number of unadjusted rows in the row group size_type source_index; // file index. - size_t compressed_size; // compressed size of the row group - size_t max_leaf_values; // maximum number of leaf values in the row group // Optional metadata pulled from the column and offset indexes, if present. std::optional> column_chunks; @@ -82,6 +80,20 @@ struct row_group_info { [[nodiscard]] bool has_offset_index() const { return column_chunks.has_value(); } }; +/** + * @brief Per-row-group inputs to read pass partitioning. + * + * Unlike `row_group_info`, the sizes here are scoped to the columns actually selected for + * reading, so they model the memory a pass will really occupy rather than the size of the + * entire row group. + */ +struct row_group_pass_size_info { + size_t start_row; // global start row of this row group + size_t num_rows; // unadjusted number of rows in this row group + size_t compressed_size; // compressed size of the selected columns in this row group + size_t max_leaf_values; // maximum number of leaf values over the selected columns +}; + /** * @brief Translates Parquet datatype to cuDF type enum */ @@ -598,15 +610,14 @@ class aggregate_reader_metadata { [[nodiscard]] std::vector get_pandas_index_names() const; /** - * @brief Computes the compressed and total size, the number of rows, and the maximum number of - * leaf values in the specified row group + * @brief Computes the compressed size, the number of rows, and the maximum number of leaf values + * in the specified row group, over all of its columns * - * @param row_group The row group + * @param rg The row group * - * @return A tuple of row group compressed size, total size, number of rows, and maximum leaf - * values + * @return A tuple of row group compressed size, number of rows, and maximum leaf values */ - [[nodiscard]] std::tuple get_row_group_properties( + [[nodiscard]] std::tuple get_row_group_properties( RowGroup const& rg) const; /** diff --git a/cpp/tests/io/parquet_chunked_reader_test.cu b/cpp/tests/io/parquet_chunked_reader_test.cu index c131a9983e01..d2ce7dc3a9c9 100644 --- a/cpp/tests/io/parquet_chunked_reader_test.cu +++ b/cpp/tests/io/parquet_chunked_reader_test.cu @@ -1338,6 +1338,85 @@ TEST_F(ParquetChunkedReaderInputLimitConstrainedTest, MixedColumns) struct ParquetChunkedReaderInputLimitTest : public cudf::test::BaseFixture {}; +namespace { +// Reads a file through the chunked reader, optionally projecting a subset of columns. +auto chunked_read_projected(std::string const& filepath, + std::vector const& columns, + std::size_t output_limit, + std::size_t input_limit) +{ + auto builder = cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}); + if (not columns.empty()) { builder.columns(columns); } + auto const read_opts = builder.build(); + auto reader = cudf::io::chunked_parquet_reader(output_limit, input_limit, read_opts); + + auto num_chunks = 0; + auto out_tables = std::vector>{}; + do { + auto chunk = reader.read_chunk(); + ++num_chunks; + out_tables.emplace_back(std::move(chunk.tbl)); + } while (reader.has_next()); + + auto out_tviews = std::vector{}; + for (auto const& tbl : out_tables) { + out_tviews.emplace_back(tbl->view()); + } + return std::pair(cudf::concatenate(out_tviews), num_chunks); +} +} // namespace + +// Pass construction must size each row group by the columns actually being read. Projecting a +// single column out of a wide file should therefore need far fewer passes - and so produce far +// fewer output chunks - than reading every column under the same `pass_read_limit`. +TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsShrinkPasses) +{ + constexpr int num_columns = 16; + constexpr int num_rows = 500'000; + + auto const filepath = temp_env->get_temp_filepath("projected_columns_shrink_passes.parquet"); + + auto const iter = cuda::counting_iterator{0}; + auto columns = std::vector>{}; + auto column_names = std::vector{}; + for (int c = 0; c < num_columns; c++) { + auto col = cudf::test::fixed_width_column_wrapper(iter, iter + num_rows); + columns.emplace_back(col.release()); + column_names.emplace_back("col_" + std::to_string(c)); + } + auto const input_table = cudf::table(std::move(columns)); + + auto metadata = cudf::io::table_input_metadata{input_table.view()}; + for (int c = 0; c < num_columns; c++) { + metadata.column_metadata[c].set_name(column_names[c]); + } + + // Many small row groups so that the pass builder has boundaries to choose between, and no + // dictionary encoding so that column chunk sizes are predictable. + cudf::io::write_parquet( + cudf::io::parquet_writer_options::builder(cudf::io::sink_info{filepath}, input_table.view()) + .metadata(std::move(metadata)) + .compression(cudf::io::compression_type::NONE) + .dictionary_policy(cudf::io::dictionary_policy::NEVER) + .row_group_size_rows(20'000) + .build()); + + // Sized so that a pass holds several row groups' worth of one column, but well under a single + // row group's worth of all 16 columns. + constexpr std::size_t pass_read_limit = 32ul * 1024 * 1024; + + auto const [all_cols, all_cols_chunks] = chunked_read_projected(filepath, {}, 0, pass_read_limit); + auto const [one_col, one_col_chunks] = + chunked_read_projected(filepath, {column_names[0]}, 0, pass_read_limit); + + // The projected read must not be split more finely than the full read. + EXPECT_LT(one_col_chunks, all_cols_chunks); + + // ...and it must still return the right data. + CUDF_TEST_EXPECT_TABLES_EQUAL(all_cols->view(), input_table.view()); + CUDF_TEST_EXPECT_TABLES_EQUAL(one_col->view(), cudf::table_view{{input_table.view().column(0)}}); +} + namespace { struct offset_gen { int const group_size; From ae043864071679542ce8df9befff2957a3da48bc Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Fri, 24 Jul 2026 18:10:59 -0700 Subject: [PATCH 2/9] Make hybrid scan pass planning column-selection aware `hybrid_scan_reader::construct_row_group_passes()` estimated each row group's compressed size over every column in the row group, which its own documentation acknowledged as a deliberately conservative approximation. That approximation is worst exactly where hybrid scan is used: a filter pass typically touches one or two columns out of a wide schema, so the planner splits into roughly `total_columns / selected_columns` times more passes than `pass_read_limit` requires. Give the planning API the reader options so it can resolve the actual column selection, and add `aggregate_reader_metadata::get_row_group_pass_size_info()`, which sums `total_compressed_size` and takes the max `num_values` over just the selected leaf column chunks. Schema indices are mapped per source with `map_schema_index()` so mismatched-schema multi-source reads stay correct. Sizing uses `ALL_COLUMNS` - the union of projected payload columns and filter columns - because this call cannot know whether the caller will use the two-stage filter/payload flow or `setup_chunking_for_all_columns()`, and the latter needs the full footprint. `select_columns()` memoizes via the `_is_*_columns_selected` flags and `reset_internal_state()` does not clear them, so planning calls `reset_column_selection()` afterwards. Without it, a later `setup_chunking_for_all_columns()` would short-circuit selection and never rebuild `_output_buffers` / `_output_buffers_template`. `get_row_group_properties()` had no remaining callers and is removed. The JNI passes the wrapper's current options, which already track any filter installed via `setFilter`, so no Java signature change is needed. Co-Authored-By: Claude Opus 5 --- .../cudf/io/experimental/hybrid_scan.hpp | 11 ++- .../io/experimental/hybrid_scan_multifile.hpp | 8 +- .../io/parquet/experimental/hybrid_scan.cpp | 7 +- .../parquet/experimental/hybrid_scan_impl.cpp | 46 +++++++--- .../parquet/experimental/hybrid_scan_impl.hpp | 2 + .../experimental/hybrid_scan_multifile.cpp | 5 +- .../io/parquet/reader_impl_chunking_utils.cuh | 3 +- cpp/src/io/parquet/reader_impl_helpers.cpp | 42 ++++----- cpp/src/io/parquet/reader_impl_helpers.hpp | 19 ++-- .../experimental/hybrid_scan_filters_test.cpp | 86 ++++++++++++++++++- .../hybrid_scan_multifile_filters_test.cpp | 18 ++-- .../io/experimental/hybrid_scan_test.cpp | 11 ++- cpp/tests/io/parquet_chunked_reader_test.cu | 18 ++-- .../src/HybridScanReaderJniMaterialize.cpp | 5 +- .../pylibcudf/io/experimental/hybrid_scan.pyi | 1 + .../pylibcudf/io/experimental/hybrid_scan.pyx | 7 +- .../pylibcudf/libcudf/io/hybrid_scan.pxd | 1 + .../tests/io/test_experimental_hybrid_scan.py | 8 +- 18 files changed, 225 insertions(+), 73 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan.hpp b/cpp/include/cudf/io/experimental/hybrid_scan.hpp index 6a3b2059b55d..7a2595f2f375 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan.hpp @@ -667,19 +667,24 @@ class hybrid_scan_reader { * * 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. + * size is estimated over only the columns selected for reading by @p options - the union of the + * projected payload columns and the filter columns. The union is used because this call cannot + * know whether the caller will use the two-stage filter/payload materialization flow or read all + * columns in one go, and the latter needs the full footprint. * * @throws std::invalid_argument if no row group indices in the input * * @param row_group_indices Input row group indices + * @param options Reader options used to determine the columns that will be read * @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 */ [[nodiscard]] std::vector> construct_row_group_passes( - std::span row_group_indices, std::size_t pass_read_limit) const; + std::span row_group_indices, + parquet_reader_options const& options, + std::size_t pass_read_limit) const; /** * @brief Check if there is any parquet data left to read for the current setup diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index 1be7365c53f4..cad5f434cd34 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -432,12 +432,15 @@ class hybrid_scan_multifile { * * 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. + * size is estimated over only the columns selected for reading by @p options - the union of the + * projected payload columns and the filter columns. The union is used because this call cannot + * know whether the caller will use the two-stage filter/payload materialization flow or read all + * columns in one go, and the latter needs the full footprint. * * @throws std::invalid_argument if no row group indices in the input * * @param row_group_indices Span of vectors of input row group indices, one per source + * @param options Reader options used to determine the columns that will be read * @param pass_read_limit Memory limit to read and decompress row group data, `0` if there is * no limit (single pass) * @@ -445,6 +448,7 @@ class hybrid_scan_multifile { */ [[nodiscard]] std::vector>> construct_row_group_passes( cudf::host_span const> row_group_indices, + parquet_reader_options const& options, std::size_t pass_read_limit) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 36bc5de06dc4..8c8759826fcc 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -368,7 +368,9 @@ table_with_metadata hybrid_scan_reader::materialize_all_columns_chunk() const } std::vector> hybrid_scan_reader::construct_row_group_passes( - std::span row_group_indices, std::size_t pass_read_limit) const + std::span row_group_indices, + parquet_reader_options const& options, + std::size_t pass_read_limit) const { CUDF_FUNC_RANGE(); @@ -381,7 +383,8 @@ std::vector> hybrid_scan_reader::construct_row_grou 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) + ->construct_row_group_passes( + input_row_group_indices, total_row_groups, options, pass_read_limit) .first; } diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 30a144db1bc4..7e8aae56b94b 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -781,6 +781,7 @@ std::pair>, std::vector const> row_group_indices, std::size_t total_row_groups, + parquet_reader_options const& options, std::size_t pass_read_limit) const { CUDF_EXPECTS( @@ -800,6 +801,34 @@ hybrid_scan_reader_impl::construct_row_group_passes( CUDF_EXPECTS( pass_read_limit > 0, "Pass read limit must be greater than 0", std::invalid_argument); + // Resolve the columns that will be read. This mirrors the `ALL_COLUMNS` case of + // `select_columns()` but resolves into locals rather than into the reader's own selection state: + // this is a planning call that may be made at any point, including while a chunked + // materialization is in flight, and `materialize_*_chunk()` do not re-select columns. + // + // `ALL_COLUMNS` - the union of the projected payload columns and the filter columns - is the + // right footprint because this call cannot know whether the caller will use the two-stage + // filter/payload flow or `setup_chunking_for_all_columns()`, and the latter needs the full + // footprint. + auto const selection_options = make_column_selection_options(options); + auto const select_column_names = get_column_projection(options); + auto filter_only_columns_names = std::optional>{}; + if (options.get_filter().has_value() and select_column_names.has_value()) { + filter_only_columns_names = parquet::detail::get_column_names_in_expression( + options.get_filter(), *select_column_names, options, _extended_metadata->get_schema_tree()); + } + // Note: this also populates the metadata's schema index maps, which `map_schema_index()` below + // relies on for sources with mismatched schemas. Doing so is idempotent. + auto const selected_columns = std::get<0>( + _metadata->select_columns(select_column_names, filter_only_columns_names, selection_options)); + + auto selected_schema_indices = std::vector{}; + selected_schema_indices.reserve(selected_columns.size()); + std::transform(selected_columns.begin(), + selected_columns.end(), + std::back_inserter(selected_schema_indices), + [](auto const& col) { return col.schema_idx; }); + // Row group (index, source index) pairs, flattened across sources, parallel to `row_group_sizes` auto row_group_ids = std::vector>{}; auto row_group_sizes = std::vector{}; @@ -813,14 +842,11 @@ hybrid_scan_reader_impl::construct_row_group_passes( for (auto const rg_index : row_group_indices[source_index]) { auto const& row_group = _extended_metadata->get_row_group(rg_index, source_index); - auto const [compressed_size, num_rows, max_leaf_values] = - _extended_metadata->get_row_group_properties(row_group); + auto const rg_size = _extended_metadata->get_row_group_pass_size_info( + row_group, start_row, source_index, selected_schema_indices); + start_row += rg_size.num_rows; row_group_ids.emplace_back(rg_index, source_index); - row_group_sizes.push_back({.start_row = start_row, - .num_rows = num_rows, - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}); - start_row += num_rows; + row_group_sizes.push_back(rg_size); } }); @@ -846,10 +872,10 @@ hybrid_scan_reader_impl::construct_row_group_passes( pass.reserve(end - start); std::for_each(row_group_ids.begin() + start, row_group_ids.begin() + end, - [&](auto const& [rg_index, source_index]) { - pass.emplace_back(rg_index); + [&](auto const& row_group_id) { + pass.emplace_back(row_group_id.first); if (has_multiple_sources) { - row_group_source_map.emplace_back(source_index); + row_group_source_map.emplace_back(row_group_id.second); } }); return pass; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 3dd241af35ff..50c950fda363 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -292,6 +292,7 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { * * @param row_group_indices Span of vectors of input row group indices, one per source * @param total_row_groups Total number of row groups across all sources + * @param options Reader options, used to determine which columns will be read * @param pass_read_limit Memory limit to read and decompress row * group data * @@ -301,6 +302,7 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { [[nodiscard]] std::pair>, std::vector> construct_row_group_passes(cudf::host_span const> row_group_indices, std::size_t total_row_groups, + parquet_reader_options const& options, std::size_t pass_read_limit) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index 38c355a651d7..7e8cc53cfe88 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -255,6 +255,7 @@ bool hybrid_scan_multifile::has_next_table_chunk() const { return _impl->has_nex std::vector>> hybrid_scan_multifile::construct_row_group_passes( cudf::host_span const> row_group_indices, + parquet_reader_options const& options, std::size_t pass_read_limit) const { CUDF_FUNC_RANGE(); @@ -267,8 +268,8 @@ std::vector>> hybrid_scan_multifile::construc 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); + auto [passes, source_map] = _impl->construct_row_group_passes( + row_group_indices, total_row_groups, options, pass_read_limit); if (pass_read_limit == 0) { return {passes}; } diff --git a/cpp/src/io/parquet/reader_impl_chunking_utils.cuh b/cpp/src/io/parquet/reader_impl_chunking_utils.cuh index b261507ff95f..1d708a3a0f38 100644 --- a/cpp/src/io/parquet/reader_impl_chunking_utils.cuh +++ b/cpp/src/io/parquet/reader_impl_chunking_utils.cuh @@ -788,7 +788,8 @@ struct row_group_pass_data { * and row count thresholds. * * The sizes in @p row_group_sizes are expected to be scoped to the columns selected for reading; - * see `make_row_group_pass_size_info()`. + * see `make_row_group_pass_size_info()` and + * `aggregate_reader_metadata::get_row_group_pass_size_info()`. * * @param row_group_sizes Span of per-row-group pass sizing properties * @param comp_read_limit Maximum compressed bytes per pass diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 4802deabc653..64ab4169762d 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -1384,27 +1384,29 @@ std::vector aggregate_reader_metadata::get_pandas_index_names() con return names; } -std::tuple aggregate_reader_metadata::get_row_group_properties( - RowGroup const& row_group) const +row_group_pass_size_info aggregate_reader_metadata::get_row_group_pass_size_info( + RowGroup const& row_group, + size_t start_row, + size_type source_index, + host_span selected_schema_indices) const { - auto const compressed_size = std::transform_reduce( - row_group.columns.cbegin(), - row_group.columns.cend(), - size_t{0}, - std::plus<>(), - [](auto const& colchunk) { return colchunk.meta_data.total_compressed_size; }); - - size_t const max_leaf_values = - row_group.columns.empty() - ? 0 - : std::max_element(row_group.columns.cbegin(), - row_group.columns.cend(), - [](auto const& a, auto const& b) { - return a.meta_data.num_values < b.meta_data.num_values; - }) - ->meta_data.num_values; - - return {compressed_size, static_cast(row_group.num_rows), max_leaf_values}; + size_t compressed_size = 0; + size_t max_leaf_values = 0; + + // Only the selected leaf columns are read into device memory, so only their column chunks + // contribute to the memory footprint of a pass. + for (auto const schema_idx : selected_schema_indices) { + auto const mapped_schema_idx = map_schema_index(schema_idx, source_index); + auto const& col_meta = + row_group.columns[find_colchunk_iter_offset(row_group, mapped_schema_idx)].meta_data; + compressed_size += col_meta.total_compressed_size; + max_leaf_values = std::max(max_leaf_values, static_cast(col_meta.num_values)); + } + + return {.start_row = start_row, + .num_rows = static_cast(row_group.num_rows), + .compressed_size = compressed_size, + .max_leaf_values = max_leaf_values}; } std::tuple get_pandas_index_names() const; /** - * @brief Computes the compressed size, the number of rows, and the maximum number of leaf values - * in the specified row group, over all of its columns + * @brief Computes the pass sizing properties of a row group, scoped to the selected columns + * + * Only the column chunks belonging to @p selected_schema_indices contribute to the returned + * compressed size and maximum leaf value count. This yields a far tighter memory estimate than + * sizing the row group over its entire schema when only a subset of a wide schema is being read. * * @param rg The row group + * @param start_row Global start row of this row group + * @param source_index Index of the data source this row group belongs to + * @param selected_schema_indices Schema indices of the selected leaf columns * - * @return A tuple of row group compressed size, number of rows, and maximum leaf values + * @return Pass sizing properties of the row group */ - [[nodiscard]] std::tuple get_row_group_properties( - RowGroup const& rg) const; + [[nodiscard]] row_group_pass_size_info get_row_group_pass_size_info( + RowGroup const& rg, + size_t start_row, + size_type source_index, + host_span selected_schema_indices) const; /** * @brief Filters the row groups using stats and bloom filters based on predicate filter diff --git a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp index 95a5677b4372..5748aa948f47 100644 --- a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp @@ -13,6 +13,7 @@ #include #include #include +#include #include #include #include @@ -1621,6 +1622,85 @@ TYPED_TEST(RowGroupFilteringWithDictTest, FilterManyLiteralsTyped) } } +// Pass construction must size each row group by the columns the reader options actually select. +// Projecting a single column out of a wide file should therefore need strictly fewer passes than +// reading every column under the same `pass_read_limit`. +TEST_F(HybridScanFiltersTest, RowGroupPassesUseSelectedColumnsOnly) +{ + auto constexpr num_columns = 8; + auto constexpr num_rg = 10; + auto constexpr rows_per_rg = 5'000; + + auto values = cuda::counting_iterator(0); + auto columns = std::vector>{}; + auto column_names = std::vector{}; + for (int c = 0; c < num_columns; c++) { + cudf::test::fixed_width_column_wrapper col(values, values + rows_per_rg); + columns.emplace_back(col.release()); + column_names.emplace_back("col_" + std::to_string(c)); + } + auto const chunk_table = cudf::table(std::move(columns)); + + auto metadata = cudf::io::table_input_metadata{chunk_table.view()}; + for (int c = 0; c < num_columns; c++) { + metadata.column_metadata[c].set_name(column_names[c]); + } + + auto const parquet_filepath = + temp_env->get_temp_filepath("RowGroupPassesSelectedColumns.parquet"); + { + // Uncompressed and undictionaried so that column chunk sizes are predictable + auto opts = + cudf::io::chunked_parquet_writer_options::builder(cudf::io::sink_info{parquet_filepath}) + .metadata(std::move(metadata)) + .compression(cudf::io::compression_type::NONE) + .dictionary_policy(cudf::io::dictionary_policy::NEVER) + .build(); + auto writer = cudf::io::chunked_parquet_writer(opts); + for (int i = 0; i < num_rg; ++i) { + writer.write(chunk_table.view()); + } + writer.close(); + } + + auto const all_columns_options = cudf::io::parquet_reader_options::builder().build(); + auto const one_column_options = + cudf::io::parquet_reader_options::builder().column_names({column_names[0]}).build(); + + auto datasource = cudf::io::datasource::create(parquet_filepath); + auto const footer_buffer = cudf::io::parquet::fetch_footer_to_host(*datasource); + auto reader = std::make_unique( + *footer_buffer, all_columns_options); + + auto const all_row_groups = reader->all_row_groups(all_columns_options); + EXPECT_EQ(static_cast(all_row_groups.size()), num_rg); + + // Each column chunk is 5'000 * sizeof(int32_t) == 20'000 bytes (plus a small page header), so a + // row group is ~160'000 bytes over all 8 columns. The pass budget is + // `pass_read_limit * input_limit_compression_reserve` == 60'000 bytes, which one row group of + // all columns already exceeds, but which fits two row groups of a single column. + auto constexpr pass_read_limit = std::size_t{200'000}; + + auto const all_column_passes = + reader->construct_row_group_passes(all_row_groups, all_columns_options, pass_read_limit); + auto const one_column_passes = + reader->construct_row_group_passes(all_row_groups, one_column_options, pass_read_limit); + + // One row group per pass over all columns; two row groups per pass for a single column. + EXPECT_EQ(all_column_passes.size(), static_cast(num_rg)); + EXPECT_EQ(one_column_passes.size(), static_cast(num_rg) / 2); + + // Both partitions must still cover every row group exactly once, in order + for (auto const& passes : {all_column_passes, one_column_passes}) { + std::vector flattened; + for (auto const& pass : passes) { + EXPECT_GT(pass.size(), 0); + flattened.insert(flattened.end(), pass.begin(), pass.end()); + } + EXPECT_EQ(flattened, all_row_groups); + } +} + TEST_F(HybridScanFiltersTest, RowGroupPasses) { auto constexpr num_rg = 10; @@ -1656,14 +1736,14 @@ TEST_F(HybridScanFiltersTest, RowGroupPasses) // No pass read limit. All row groups in a single pass { - auto passes = reader->construct_row_group_passes(all_row_groups, 0); + auto passes = reader->construct_row_group_passes(all_row_groups, options, 0); EXPECT_EQ(passes.size(), 1); EXPECT_EQ(passes.front(), all_row_groups); } // Small pass limit would result in each row group in its own pass { - auto passes = reader->construct_row_group_passes(all_row_groups, 1); + auto passes = reader->construct_row_group_passes(all_row_groups, options, 1); EXPECT_EQ(passes.size(), all_row_groups.size()); auto zipped = cuda::make_zip_iterator(passes.begin(), all_row_groups.begin()); std::for_each(zipped, zipped + passes.size(), [&](auto const& iter) { @@ -1676,7 +1756,7 @@ TEST_F(HybridScanFiltersTest, RowGroupPasses) // All passes should cover all row groups and be consecutive { - auto passes = reader->construct_row_group_passes(all_row_groups, 1'024); + auto passes = reader->construct_row_group_passes(all_row_groups, options, 1'024); std::vector flattened; for (auto const& pass : passes) { EXPECT_GT(pass.size(), 0); 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 baf05e6188b9..372fca674897 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -231,7 +231,7 @@ TEST_F(HybridScanMultifileFiltersTest, EmptySource) 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); + auto const passes = reader->construct_row_group_passes(all_rgs, options, 1); ASSERT_EQ(passes.size(), all_rgs.front().size()); for (auto const& pass : passes) { ASSERT_EQ(pass.size(), num_sources); @@ -270,22 +270,22 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) { auto invalid_rgs = all_rgs; invalid_rgs.pop_back(); - EXPECT_THROW(static_cast(reader->construct_row_group_passes(invalid_rgs, 0)), + EXPECT_THROW(static_cast(reader->construct_row_group_passes(invalid_rgs, options, 0)), 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)), + EXPECT_THROW(static_cast(reader->construct_row_group_passes(empty_rgs, options, 0)), std::invalid_argument); - EXPECT_THROW(static_cast(reader->construct_row_group_passes(empty_rgs, 1)), + EXPECT_THROW(static_cast(reader->construct_row_group_passes(empty_rgs, options, 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); + auto const passes = reader->construct_row_group_passes(all_rgs, options, 0); ASSERT_EQ(passes.size(), 1); ASSERT_EQ(passes.front().size(), num_sources); EXPECT_EQ(passes.front(), all_rgs); @@ -293,7 +293,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) // Small pass read limit => each row group in its own pass { - auto const passes = reader->construct_row_group_passes(all_rgs, 1); + auto const passes = reader->construct_row_group_passes(all_rgs, options, 1); ASSERT_EQ(passes.size(), num_sources * all_rgs.front().size()); for (auto const& pass : passes) { ASSERT_EQ(pass.size(), num_sources); @@ -307,7 +307,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) // Large pass read limit => multiple passes { - auto const passes = reader->construct_row_group_passes(all_rgs, 10'000); + auto const passes = reader->construct_row_group_passes(all_rgs, options, 10'000); ASSERT_GT(passes.size(), 1); auto const pass_num_row_groups = [](auto const& pass) { return std::accumulate( @@ -366,9 +366,9 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPassesSingleSourceParity) 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); + multifile_reader->construct_row_group_passes(all_rgs, options, pass_read_limit); auto const single_file_passes = - single_file_reader->construct_row_group_passes(all_rgs.front(), pass_read_limit); + single_file_reader->construct_row_group_passes(all_rgs.front(), options, pass_read_limit); auto projected_passes = std::vector>{}; projected_passes.reserve(multifile_passes.size()); diff --git a/cpp/tests/io/experimental/hybrid_scan_test.cpp b/cpp/tests/io/experimental/hybrid_scan_test.cpp index edf43fc6bdff..aa56fd2f1ea4 100644 --- a/cpp/tests/io/experimental/hybrid_scan_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_test.cpp @@ -1206,7 +1206,8 @@ TEST_F(HybridScanTest, RowGroupPassesMatchesChunkedReader) *footer_buffer, options); auto const all_row_groups = reader->all_row_groups(options); - auto const passes = reader->construct_row_group_passes(all_row_groups, pass_read_limit); + auto const passes = + reader->construct_row_group_passes(all_row_groups, options, pass_read_limit); for (auto const& pass_row_groups : passes) { auto const chunk_byte_ranges = @@ -1218,7 +1219,15 @@ TEST_F(HybridScanTest, RowGroupPassesMatchesChunkedReader) reader->setup_chunking_for_all_columns( 0, pass_read_limit, pass_row_groups, col_data, options, stream, mr); + // Pass planning must not disturb an in-flight chunked materialization. + // `materialize_all_columns_chunk()` does not re-select columns, so a planner call that + // mutated the reader's column selection would silently swap this loop onto a single column. + auto const probe_options = + cudf::io::parquet_reader_options::builder().column_indices({0}).build(); + while (reader->has_next_table_chunk()) { + static_cast( + reader->construct_row_group_passes(pass_row_groups, probe_options, pass_read_limit)); auto chunk = reader->materialize_all_columns_chunk(); hybrid_scan_tables.push_back(std::move(chunk.tbl)); } diff --git a/cpp/tests/io/parquet_chunked_reader_test.cu b/cpp/tests/io/parquet_chunked_reader_test.cu index d2ce7dc3a9c9..9c4e762e1033 100644 --- a/cpp/tests/io/parquet_chunked_reader_test.cu +++ b/cpp/tests/io/parquet_chunked_reader_test.cu @@ -1346,7 +1346,7 @@ auto chunked_read_projected(std::string const& filepath, std::size_t input_limit) { auto builder = cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}); - if (not columns.empty()) { builder.columns(columns); } + if (not columns.empty()) { builder.column_names(columns); } auto const read_opts = builder.build(); auto reader = cudf::io::chunked_parquet_reader(output_limit, input_limit, read_opts); @@ -1372,15 +1372,15 @@ auto chunked_read_projected(std::string const& filepath, TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsShrinkPasses) { constexpr int num_columns = 16; - constexpr int num_rows = 500'000; + constexpr int num_rows = 125'000; auto const filepath = temp_env->get_temp_filepath("projected_columns_shrink_passes.parquet"); - auto const iter = cuda::counting_iterator{0}; + auto const iter = cuda::counting_iterator{0}; auto columns = std::vector>{}; auto column_names = std::vector{}; for (int c = 0; c < num_columns; c++) { - auto col = cudf::test::fixed_width_column_wrapper(iter, iter + num_rows); + auto col = cudf::test::fixed_width_column_wrapper(iter, iter + num_rows); columns.emplace_back(col.release()); column_names.emplace_back("col_" + std::to_string(c)); } @@ -1398,12 +1398,14 @@ TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsShrinkPasses) .metadata(std::move(metadata)) .compression(cudf::io::compression_type::NONE) .dictionary_policy(cudf::io::dictionary_policy::NEVER) - .row_group_size_rows(20'000) + .row_group_size_rows(5'000) .build()); - // Sized so that a pass holds several row groups' worth of one column, but well under a single - // row group's worth of all 16 columns. - constexpr std::size_t pass_read_limit = 32ul * 1024 * 1024; + // 25 row groups, each column chunk 5'000 * sizeof(int32_t) == 20'000 bytes, so a row group is + // ~320'000 bytes over all 16 columns. The pass budget is + // `pass_read_limit * input_limit_compression_reserve` == ~400'000 bytes: one row group per pass + // for the full read, but twenty row groups per pass for a single column. + constexpr std::size_t pass_read_limit = 1'333'334; auto const [all_cols, all_cols_chunks] = chunked_read_projected(filepath, {}, 0, pass_read_limit); auto const [one_col, one_col_chunks] = diff --git a/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp b/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp index f6f671ae4e7d..d41398610adb 100644 --- a/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp +++ b/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp @@ -364,8 +364,9 @@ JNIEXPORT jobjectArray JNICALL Java_ai_rapids_cudf_HybridScanReader_constructRow auto const pass_limit = checked_size_t(env, pass_read_limit, "pass_read_limit"); auto* wrapper = reinterpret_cast(handle); auto holder = make_row_group_span(env, j_row_groups); - auto passes = wrapper->reader->construct_row_group_passes(holder.span(), pass_limit); - jclass int_array_cls = env->FindClass("[I"); + auto passes = + wrapper->reader->construct_row_group_passes(holder.span(), wrapper->options, pass_limit); + jclass int_array_cls = env->FindClass("[I"); if (int_array_cls == nullptr) { return nullptr; } auto outer = env->NewObjectArray(passes.size(), int_array_cls, nullptr); if (outer == nullptr) { return nullptr; } diff --git a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi index f95dc8b054d3..11edfaa0ef0a 100644 --- a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi +++ b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi @@ -139,6 +139,7 @@ class HybridScanReader: def construct_row_group_passes( self, row_group_indices: list[int], + options: ParquetReaderOptions, pass_read_limit: int, ) -> list[list[int]]: ... def has_next_table_chunk(self) -> bool: ... diff --git a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx index 664eb489428c..2bde59ba088a 100644 --- a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx +++ b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx @@ -789,6 +789,7 @@ cdef class HybridScanReader: def construct_row_group_passes( self, list row_group_indices, + ParquetReaderOptions options, size_t pass_read_limit, ): """Partition row groups into passes such that the GPU memory required to @@ -796,12 +797,15 @@ cdef class HybridScanReader: Note that ``pass_read_limit`` is a hint, not an absolute limit. i.e. if a row group cannot fit within the limit, it will still constitute a valid - pass. + pass. Row group sizes are estimated over only the columns selected for + reading by ``options``. Parameters ---------- row_group_indices : list[int] Input row group indices + options : ParquetReaderOptions + Parquet reader options, used to determine the columns that will be read pass_read_limit : int Limit on the amount of memory used for reading and decompressing data or 0 if there is no limit. @@ -821,6 +825,7 @@ cdef class HybridScanReader: std_span[const_size_type]( indices_vec.data(), indices_vec.size() ), + options.c_obj, pass_read_limit ) diff --git a/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd b/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd index 36201d545de6..4be97f006f94 100644 --- a/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd +++ b/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd @@ -170,6 +170,7 @@ cdef extern from "cudf/io/experimental/hybrid_scan.hpp" \ vector[vector[size_type]] construct_row_group_passes( std_span[const_size_type] row_group_indices, + const parquet_reader_options& options, size_t pass_read_limit, ) except +libcudf_exception_handler diff --git a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py index 74f467f16193..f0a874ab09ed 100644 --- a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py +++ b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py @@ -680,21 +680,21 @@ def test_hybrid_scan_construct_row_group_passes( # zero pass read limit => single pass with all row groups pass_read_limit = 0 passes = simple_hybrid_scan_reader.construct_row_group_passes( - all_row_groups, pass_read_limit + all_row_groups, simple_parquet_options, pass_read_limit ) assert passes == [all_row_groups] # small pass read limit => each row group in its own pass pass_read_limit = 1 passes = simple_hybrid_scan_reader.construct_row_group_passes( - all_row_groups, pass_read_limit + all_row_groups, simple_parquet_options, pass_read_limit ) assert passes == [[rg] for rg in all_row_groups] # Passes should flatten to all row groups pass_read_limit = 1024 passes = simple_hybrid_scan_reader.construct_row_group_passes( - all_row_groups, pass_read_limit + all_row_groups, simple_parquet_options, pass_read_limit ) assert [rg for p in passes for rg in p] == all_row_groups assert all(passes) @@ -704,7 +704,7 @@ def test_hybrid_scan_construct_row_group_passes( ValueError, match="Empty input row group indices encountered" ): simple_hybrid_scan_reader.construct_row_group_passes( - [], pass_read_limit + [], simple_parquet_options, pass_read_limit ) From 95153726f23ded9c5030e4ef0bb1a94810916ac2 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Mon, 27 Jul 2026 23:27:35 +0000 Subject: [PATCH 3/9] Improve Parquet reader pass construction Size passes from selected columns, simplify row-group sizing, and retain hybrid scan compatibility. --- .../cudf/io/experimental/hybrid_scan.hpp | 11 +-- .../io/experimental/hybrid_scan_multifile.hpp | 8 +- .../io/parquet/experimental/hybrid_scan.cpp | 7 +- .../parquet/experimental/hybrid_scan_impl.cpp | 42 ++------- .../parquet/experimental/hybrid_scan_impl.hpp | 2 - .../experimental/hybrid_scan_multifile.cpp | 5 +- cpp/src/io/parquet/reader_impl_chunking.cu | 2 +- .../io/parquet/reader_impl_chunking_utils.cu | 59 +++++++------ .../io/parquet/reader_impl_chunking_utils.cuh | 27 ++---- cpp/src/io/parquet/reader_impl_helpers.cpp | 42 +++++---- cpp/src/io/parquet/reader_impl_helpers.hpp | 42 +++------ .../experimental/hybrid_scan_filters_test.cpp | 86 +------------------ .../hybrid_scan_multifile_filters_test.cpp | 18 ++-- .../io/experimental/hybrid_scan_test.cpp | 11 +-- .../src/HybridScanReaderJniMaterialize.cpp | 5 +- .../pylibcudf/io/experimental/hybrid_scan.pyi | 1 - .../pylibcudf/io/experimental/hybrid_scan.pyx | 7 +- .../pylibcudf/libcudf/io/hybrid_scan.pxd | 1 - .../tests/io/test_experimental_hybrid_scan.py | 8 +- 19 files changed, 111 insertions(+), 273 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan.hpp b/cpp/include/cudf/io/experimental/hybrid_scan.hpp index 7a2595f2f375..6a3b2059b55d 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan.hpp @@ -667,24 +667,19 @@ class hybrid_scan_reader { * * 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 only the columns selected for reading by @p options - the union of the - * projected payload columns and the filter columns. The union is used because this call cannot - * know whether the caller will use the two-stage filter/payload materialization flow or read all - * columns in one go, and the latter needs the full footprint. + * size is estimated over all columns in each row group (not just the columns selected for + * reading), for conservative estimates. * * @throws std::invalid_argument if no row group indices in the input * * @param row_group_indices Input row group indices - * @param options Reader options used to determine the columns that will be read * @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 */ [[nodiscard]] std::vector> construct_row_group_passes( - std::span row_group_indices, - parquet_reader_options const& options, - std::size_t pass_read_limit) const; + std::span row_group_indices, std::size_t pass_read_limit) const; /** * @brief Check if there is any parquet data left to read for the current setup diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index cad5f434cd34..1be7365c53f4 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -432,15 +432,12 @@ class hybrid_scan_multifile { * * 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 only the columns selected for reading by @p options - the union of the - * projected payload columns and the filter columns. The union is used because this call cannot - * know whether the caller will use the two-stage filter/payload materialization flow or read all - * columns in one go, and the latter needs the full footprint. + * size is estimated over all columns in each row group (not just the columns selected for + * reading), for conservative estimates. * * @throws std::invalid_argument if no row group indices in the input * * @param row_group_indices Span of vectors of input row group indices, one per source - * @param options Reader options used to determine the columns that will be read * @param pass_read_limit Memory limit to read and decompress row group data, `0` if there is * no limit (single pass) * @@ -448,7 +445,6 @@ class hybrid_scan_multifile { */ [[nodiscard]] std::vector>> construct_row_group_passes( cudf::host_span const> row_group_indices, - parquet_reader_options const& options, std::size_t pass_read_limit) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 8c8759826fcc..36bc5de06dc4 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -368,9 +368,7 @@ table_with_metadata hybrid_scan_reader::materialize_all_columns_chunk() const } std::vector> hybrid_scan_reader::construct_row_group_passes( - std::span row_group_indices, - parquet_reader_options const& options, - std::size_t pass_read_limit) const + std::span row_group_indices, std::size_t pass_read_limit) const { CUDF_FUNC_RANGE(); @@ -383,8 +381,7 @@ std::vector> hybrid_scan_reader::construct_row_grou std::vector>{{row_group_indices.begin(), row_group_indices.end()}}; return _impl - ->construct_row_group_passes( - input_row_group_indices, total_row_groups, options, pass_read_limit) + ->construct_row_group_passes(input_row_group_indices, total_row_groups, pass_read_limit) .first; } diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 7e8aae56b94b..1e3a309bb37c 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -781,7 +781,6 @@ std::pair>, std::vector const> row_group_indices, std::size_t total_row_groups, - parquet_reader_options const& options, std::size_t pass_read_limit) const { CUDF_EXPECTS( @@ -801,52 +800,23 @@ hybrid_scan_reader_impl::construct_row_group_passes( CUDF_EXPECTS( pass_read_limit > 0, "Pass read limit must be greater than 0", std::invalid_argument); - // Resolve the columns that will be read. This mirrors the `ALL_COLUMNS` case of - // `select_columns()` but resolves into locals rather than into the reader's own selection state: - // this is a planning call that may be made at any point, including while a chunked - // materialization is in flight, and `materialize_*_chunk()` do not re-select columns. - // - // `ALL_COLUMNS` - the union of the projected payload columns and the filter columns - is the - // right footprint because this call cannot know whether the caller will use the two-stage - // filter/payload flow or `setup_chunking_for_all_columns()`, and the latter needs the full - // footprint. - auto const selection_options = make_column_selection_options(options); - auto const select_column_names = get_column_projection(options); - auto filter_only_columns_names = std::optional>{}; - if (options.get_filter().has_value() and select_column_names.has_value()) { - filter_only_columns_names = parquet::detail::get_column_names_in_expression( - options.get_filter(), *select_column_names, options, _extended_metadata->get_schema_tree()); - } - // Note: this also populates the metadata's schema index maps, which `map_schema_index()` below - // relies on for sources with mismatched schemas. Doing so is idempotent. - auto const selected_columns = std::get<0>( - _metadata->select_columns(select_column_names, filter_only_columns_names, selection_options)); - - auto selected_schema_indices = std::vector{}; - selected_schema_indices.reserve(selected_columns.size()); - std::transform(selected_columns.begin(), - selected_columns.end(), - std::back_inserter(selected_schema_indices), - [](auto const& col) { return col.schema_idx; }); - - // Row group (index, source index) pairs, flattened across sources, parallel to `row_group_sizes` auto row_group_ids = std::vector>{}; - auto row_group_sizes = std::vector{}; + auto row_group_sizes = std::vector{}; row_group_ids.reserve(total_row_groups); row_group_sizes.reserve(total_row_groups); - size_t start_row = 0; std::for_each(cuda::counting_iterator(0), cuda::counting_iterator(row_group_indices.size()), [&](auto const source_index) { for (auto const rg_index : row_group_indices[source_index]) { auto const& row_group = _extended_metadata->get_row_group(rg_index, source_index); - auto const rg_size = _extended_metadata->get_row_group_pass_size_info( - row_group, start_row, source_index, selected_schema_indices); - start_row += rg_size.num_rows; + auto const [compressed_size, num_rows, max_leaf_values] = + _extended_metadata->get_row_group_properties(row_group); row_group_ids.emplace_back(rg_index, source_index); - row_group_sizes.push_back(rg_size); + row_group_sizes.push_back({.unadjusted_num_rows = num_rows, + .compressed_size = compressed_size, + .max_leaf_values = max_leaf_values}); } }); diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 50c950fda363..3dd241af35ff 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -292,7 +292,6 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { * * @param row_group_indices Span of vectors of input row group indices, one per source * @param total_row_groups Total number of row groups across all sources - * @param options Reader options, used to determine which columns will be read * @param pass_read_limit Memory limit to read and decompress row * group data * @@ -302,7 +301,6 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { [[nodiscard]] std::pair>, std::vector> construct_row_group_passes(cudf::host_span const> row_group_indices, std::size_t total_row_groups, - parquet_reader_options const& options, std::size_t pass_read_limit) const; /** diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index 7e8cc53cfe88..38c355a651d7 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -255,7 +255,6 @@ bool hybrid_scan_multifile::has_next_table_chunk() const { return _impl->has_nex std::vector>> hybrid_scan_multifile::construct_row_group_passes( cudf::host_span const> row_group_indices, - parquet_reader_options const& options, std::size_t pass_read_limit) const { CUDF_FUNC_RANGE(); @@ -268,8 +267,8 @@ std::vector>> hybrid_scan_multifile::construc 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, options, pass_read_limit); + 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}; } diff --git a/cpp/src/io/parquet/reader_impl_chunking.cu b/cpp/src/io/parquet/reader_impl_chunking.cu index 8acc978c3173..bba5f5038192 100644 --- a/cpp/src/io/parquet/reader_impl_chunking.cu +++ b/cpp/src/io/parquet/reader_impl_chunking.cu @@ -548,7 +548,7 @@ void reader_impl::compute_input_passes(read_mode mode) // input column), so this is an exact measure of the compressed bytes a pass will hold rather than // an estimate over the whole row group. auto const row_group_sizes = - make_row_group_pass_size_info(row_groups_info, _file_itm_data.chunks, _input_columns.size()); + compute_row_group_size_info(row_groups_info, _file_itm_data.chunks, _input_columns.size()); auto pass_data = compute_row_group_passes(row_group_sizes, comp_read_limit, _file_itm_data.global_skip_rows); diff --git a/cpp/src/io/parquet/reader_impl_chunking_utils.cu b/cpp/src/io/parquet/reader_impl_chunking_utils.cu index 2430a5be8f23..7b7073d05c94 100644 --- a/cpp/src/io/parquet/reader_impl_chunking_utils.cu +++ b/cpp/src/io/parquet/reader_impl_chunking_utils.cu @@ -953,40 +953,41 @@ rmm::device_uvector compute_level_decode_sizes(device_span make_row_group_pass_size_info( - cudf::host_span row_groups_info, - cudf::host_span chunks, +std::vector compute_row_group_size_info( + std::span row_groups_info, + std::span chunks, size_t num_input_columns) { CUDF_EXPECTS(chunks.size() == row_groups_info.size() * num_input_columns, "Mismatch between the number of column chunk descriptors and row groups"); - std::vector row_group_sizes; + std::vector row_group_sizes; row_group_sizes.reserve(row_groups_info.size()); - for (size_t rg_idx = 0; rg_idx < row_groups_info.size(); rg_idx++) { - auto const& rgi = row_groups_info[rg_idx]; - auto const chunks_begin = chunks.begin() + (rg_idx * num_input_columns); - auto const chunks_end = chunks_begin + num_input_columns; - size_t compressed_size = 0; - size_t max_leaf_values = 0; - std::for_each(chunks_begin, chunks_end, [&](auto const& chunk) { - compressed_size += chunk.compressed_size; - max_leaf_values = std::max(max_leaf_values, chunk.num_values); - }); - row_group_sizes.push_back({.start_row = rgi.start_row, - .num_rows = rgi.unadjusted_num_rows, - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}); - } + std::transform(cuda::counting_iterator(0), + cuda::counting_iterator(row_groups_info.size()), + std::back_inserter(row_group_sizes), + [&](auto const rg_idx) { + auto const& rg_info = row_groups_info[rg_idx]; + auto const chunks_begin = chunks.begin() + (rg_idx * num_input_columns); + auto const chunks_end = chunks_begin + num_input_columns; + size_t compressed_size = 0; + size_t max_leaf_values = 0; + std::for_each(chunks_begin, chunks_end, [&](auto const& chunk) { + compressed_size += chunk.compressed_size; + max_leaf_values = std::max(max_leaf_values, chunk.num_values); + }); + return row_group_size_info{.unadjusted_num_rows = rg_info.unadjusted_num_rows, + .compressed_size = compressed_size, + .max_leaf_values = max_leaf_values}; + }); return row_group_sizes; } -row_group_pass_data compute_row_group_passes( - cudf::host_span row_group_sizes, - std::size_t comp_read_limit, - int64_t skip_rows) +row_group_pass_data compute_row_group_passes(std::span row_group_sizes, + std::size_t comp_read_limit, + int64_t skip_rows) { auto constexpr max_rows_per_pass = static_cast(std::numeric_limits::max()); @@ -1006,9 +1007,13 @@ row_group_pass_data compute_row_group_passes( // We must use the effective size of the first row group we are reading to accurately calculate // the first non-zero `input_pass_start_row_count` unless we are reading only one row group - auto const row_group_rows = (skip_rows and row_group_sizes.size() > 1) - ? (rgi.start_row + rgi.num_rows - skip_rows) - : rgi.num_rows; + auto row_group_rows = rgi.unadjusted_num_rows; + if (row_group_sizes.size() > 1) { + CUDF_EXPECTS(std::cmp_greater_equal(rgi.unadjusted_num_rows, skip_rows), + "Row groups must contribute non-negative effective rows", + std::invalid_argument); + row_group_rows -= skip_rows; + } auto const compressed_rg_size = rgi.compressed_size; auto const row_group_leaf_values = rgi.max_leaf_values; @@ -1030,7 +1035,7 @@ row_group_pass_data compute_row_group_passes( // We always need to include at least one row group, so end the pass at the end of the current // row group if (cur_rg_start == cur_rg_index) { - CUDF_EXPECTS(std::cmp_less_equal(rgi.num_rows, max_rows_per_pass), + CUDF_EXPECTS(std::cmp_less_equal(rgi.unadjusted_num_rows, max_rows_per_pass), "Number of rows in each row group must be smaller than the column size limit"); result.pass_row_group_offsets.push_back(cur_rg_index + 1); result.pass_start_row_counts.push_back(cur_row_count + row_group_rows); diff --git a/cpp/src/io/parquet/reader_impl_chunking_utils.cuh b/cpp/src/io/parquet/reader_impl_chunking_utils.cuh index 1d708a3a0f38..49cdbfd7e6cc 100644 --- a/cpp/src/io/parquet/reader_impl_chunking_utils.cuh +++ b/cpp/src/io/parquet/reader_impl_chunking_utils.cuh @@ -787,36 +787,27 @@ struct row_group_pass_data { * @brief Partition row groups into passes based on compressed size, leaf value count, * and row count thresholds. * - * The sizes in @p row_group_sizes are expected to be scoped to the columns selected for reading; - * see `make_row_group_pass_size_info()` and - * `aggregate_reader_metadata::get_row_group_pass_size_info()`. - * - * @param row_group_sizes Span of per-row-group pass sizing properties + * @param row_group_sizes Span of row group size information * @param comp_read_limit Maximum compressed bytes per pass * @param skip_rows Number of leading rows to skip (affects the effective size of the first row * group) * @return A row_group_pass_data containing pass boundary offsets and cumulative row counts */ -row_group_pass_data compute_row_group_passes( - cudf::host_span row_group_sizes, - std::size_t comp_read_limit, - int64_t skip_rows); +row_group_pass_data compute_row_group_passes(std::span row_group_sizes, + std::size_t comp_read_limit, + int64_t skip_rows); /** - * @brief Computes per-row-group pass sizing properties from already-built column chunk descriptors - * - * @p chunks is expected to be laid out as `[row_group][input_column]`, i.e. the layout produced by - * `reader_impl::create_global_chunk_info()`. Because those descriptors only exist for the selected - * columns, the resulting sizes are exact rather than estimated. + * @brief Computes row group size information from already-built column chunk descriptors * * @param row_groups_info Span of row group metadata * @param chunks Column chunk descriptors for the selected columns of every row group * @param num_input_columns Number of selected (leaf) input columns - * @return Per-row-group pass sizing properties, one entry per row group + * @return Row group size information, one entry per row group */ -std::vector make_row_group_pass_size_info( - cudf::host_span row_groups_info, - cudf::host_span chunks, +std::vector compute_row_group_size_info( + std::span row_groups_info, + std::span chunks, size_t num_input_columns); } // namespace cudf::io::parquet::detail diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 64ab4169762d..4802deabc653 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -1384,29 +1384,27 @@ std::vector aggregate_reader_metadata::get_pandas_index_names() con return names; } -row_group_pass_size_info aggregate_reader_metadata::get_row_group_pass_size_info( - RowGroup const& row_group, - size_t start_row, - size_type source_index, - host_span selected_schema_indices) const +std::tuple aggregate_reader_metadata::get_row_group_properties( + RowGroup const& row_group) const { - size_t compressed_size = 0; - size_t max_leaf_values = 0; - - // Only the selected leaf columns are read into device memory, so only their column chunks - // contribute to the memory footprint of a pass. - for (auto const schema_idx : selected_schema_indices) { - auto const mapped_schema_idx = map_schema_index(schema_idx, source_index); - auto const& col_meta = - row_group.columns[find_colchunk_iter_offset(row_group, mapped_schema_idx)].meta_data; - compressed_size += col_meta.total_compressed_size; - max_leaf_values = std::max(max_leaf_values, static_cast(col_meta.num_values)); - } - - return {.start_row = start_row, - .num_rows = static_cast(row_group.num_rows), - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}; + auto const compressed_size = std::transform_reduce( + row_group.columns.cbegin(), + row_group.columns.cend(), + size_t{0}, + std::plus<>(), + [](auto const& colchunk) { return colchunk.meta_data.total_compressed_size; }); + + size_t const max_leaf_values = + row_group.columns.empty() + ? 0 + : std::max_element(row_group.columns.cbegin(), + row_group.columns.cend(), + [](auto const& a, auto const& b) { + return a.meta_data.num_values < b.meta_data.num_values; + }) + ->meta_data.num_values; + + return {compressed_size, static_cast(row_group.num_rows), max_leaf_values}; } std::tuple get_pandas_index_names() const; /** - * @brief Computes the pass sizing properties of a row group, scoped to the selected columns - * - * Only the column chunks belonging to @p selected_schema_indices contribute to the returned - * compressed size and maximum leaf value count. This yields a far tighter memory estimate than - * sizing the row group over its entire schema when only a subset of a wide schema is being read. + * @brief Computes the compressed size, the number of rows, and the maximum number of leaf values + * in the specified row group over all of its columns * - * @param rg The row group - * @param start_row Global start row of this row group - * @param source_index Index of the data source this row group belongs to - * @param selected_schema_indices Schema indices of the selected leaf columns + * @param row_group Input row group * - * @return Pass sizing properties of the row group + * @return A tuple of row group compressed size, number of rows, and maximum leaf values */ - [[nodiscard]] row_group_pass_size_info get_row_group_pass_size_info( - RowGroup const& rg, - size_t start_row, - size_type source_index, - host_span selected_schema_indices) const; + [[nodiscard]] std::tuple get_row_group_properties( + RowGroup const& row_group) const; /** * @brief Filters the row groups using stats and bloom filters based on predicate filter diff --git a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp index 5748aa948f47..95a5677b4372 100644 --- a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp @@ -13,7 +13,6 @@ #include #include #include -#include #include #include #include @@ -1622,85 +1621,6 @@ TYPED_TEST(RowGroupFilteringWithDictTest, FilterManyLiteralsTyped) } } -// Pass construction must size each row group by the columns the reader options actually select. -// Projecting a single column out of a wide file should therefore need strictly fewer passes than -// reading every column under the same `pass_read_limit`. -TEST_F(HybridScanFiltersTest, RowGroupPassesUseSelectedColumnsOnly) -{ - auto constexpr num_columns = 8; - auto constexpr num_rg = 10; - auto constexpr rows_per_rg = 5'000; - - auto values = cuda::counting_iterator(0); - auto columns = std::vector>{}; - auto column_names = std::vector{}; - for (int c = 0; c < num_columns; c++) { - cudf::test::fixed_width_column_wrapper col(values, values + rows_per_rg); - columns.emplace_back(col.release()); - column_names.emplace_back("col_" + std::to_string(c)); - } - auto const chunk_table = cudf::table(std::move(columns)); - - auto metadata = cudf::io::table_input_metadata{chunk_table.view()}; - for (int c = 0; c < num_columns; c++) { - metadata.column_metadata[c].set_name(column_names[c]); - } - - auto const parquet_filepath = - temp_env->get_temp_filepath("RowGroupPassesSelectedColumns.parquet"); - { - // Uncompressed and undictionaried so that column chunk sizes are predictable - auto opts = - cudf::io::chunked_parquet_writer_options::builder(cudf::io::sink_info{parquet_filepath}) - .metadata(std::move(metadata)) - .compression(cudf::io::compression_type::NONE) - .dictionary_policy(cudf::io::dictionary_policy::NEVER) - .build(); - auto writer = cudf::io::chunked_parquet_writer(opts); - for (int i = 0; i < num_rg; ++i) { - writer.write(chunk_table.view()); - } - writer.close(); - } - - auto const all_columns_options = cudf::io::parquet_reader_options::builder().build(); - auto const one_column_options = - cudf::io::parquet_reader_options::builder().column_names({column_names[0]}).build(); - - auto datasource = cudf::io::datasource::create(parquet_filepath); - auto const footer_buffer = cudf::io::parquet::fetch_footer_to_host(*datasource); - auto reader = std::make_unique( - *footer_buffer, all_columns_options); - - auto const all_row_groups = reader->all_row_groups(all_columns_options); - EXPECT_EQ(static_cast(all_row_groups.size()), num_rg); - - // Each column chunk is 5'000 * sizeof(int32_t) == 20'000 bytes (plus a small page header), so a - // row group is ~160'000 bytes over all 8 columns. The pass budget is - // `pass_read_limit * input_limit_compression_reserve` == 60'000 bytes, which one row group of - // all columns already exceeds, but which fits two row groups of a single column. - auto constexpr pass_read_limit = std::size_t{200'000}; - - auto const all_column_passes = - reader->construct_row_group_passes(all_row_groups, all_columns_options, pass_read_limit); - auto const one_column_passes = - reader->construct_row_group_passes(all_row_groups, one_column_options, pass_read_limit); - - // One row group per pass over all columns; two row groups per pass for a single column. - EXPECT_EQ(all_column_passes.size(), static_cast(num_rg)); - EXPECT_EQ(one_column_passes.size(), static_cast(num_rg) / 2); - - // Both partitions must still cover every row group exactly once, in order - for (auto const& passes : {all_column_passes, one_column_passes}) { - std::vector flattened; - for (auto const& pass : passes) { - EXPECT_GT(pass.size(), 0); - flattened.insert(flattened.end(), pass.begin(), pass.end()); - } - EXPECT_EQ(flattened, all_row_groups); - } -} - TEST_F(HybridScanFiltersTest, RowGroupPasses) { auto constexpr num_rg = 10; @@ -1736,14 +1656,14 @@ TEST_F(HybridScanFiltersTest, RowGroupPasses) // No pass read limit. All row groups in a single pass { - auto passes = reader->construct_row_group_passes(all_row_groups, options, 0); + auto passes = reader->construct_row_group_passes(all_row_groups, 0); EXPECT_EQ(passes.size(), 1); EXPECT_EQ(passes.front(), all_row_groups); } // Small pass limit would result in each row group in its own pass { - auto passes = reader->construct_row_group_passes(all_row_groups, options, 1); + auto passes = reader->construct_row_group_passes(all_row_groups, 1); EXPECT_EQ(passes.size(), all_row_groups.size()); auto zipped = cuda::make_zip_iterator(passes.begin(), all_row_groups.begin()); std::for_each(zipped, zipped + passes.size(), [&](auto const& iter) { @@ -1756,7 +1676,7 @@ TEST_F(HybridScanFiltersTest, RowGroupPasses) // All passes should cover all row groups and be consecutive { - auto passes = reader->construct_row_group_passes(all_row_groups, options, 1'024); + auto passes = reader->construct_row_group_passes(all_row_groups, 1'024); std::vector flattened; for (auto const& pass : passes) { EXPECT_GT(pass.size(), 0); 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 372fca674897..baf05e6188b9 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -231,7 +231,7 @@ TEST_F(HybridScanMultifileFiltersTest, EmptySource) 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, options, 1); + 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); @@ -270,22 +270,22 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) { auto invalid_rgs = all_rgs; invalid_rgs.pop_back(); - EXPECT_THROW(static_cast(reader->construct_row_group_passes(invalid_rgs, options, 0)), + EXPECT_THROW(static_cast(reader->construct_row_group_passes(invalid_rgs, 0)), 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, options, 0)), + 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, options, 1)), + 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, options, 0); + 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); @@ -293,7 +293,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) // Small pass read limit => each row group in its own pass { - auto const passes = reader->construct_row_group_passes(all_rgs, options, 1); + 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); @@ -307,7 +307,7 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPasses) // Large pass read limit => multiple passes { - auto const passes = reader->construct_row_group_passes(all_rgs, options, 10'000); + 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( @@ -366,9 +366,9 @@ TEST_F(HybridScanMultifileFiltersTest, RowGroupPassesSingleSourceParity) 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, options, pass_read_limit); + 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(), options, pass_read_limit); + single_file_reader->construct_row_group_passes(all_rgs.front(), pass_read_limit); auto projected_passes = std::vector>{}; projected_passes.reserve(multifile_passes.size()); diff --git a/cpp/tests/io/experimental/hybrid_scan_test.cpp b/cpp/tests/io/experimental/hybrid_scan_test.cpp index aa56fd2f1ea4..edf43fc6bdff 100644 --- a/cpp/tests/io/experimental/hybrid_scan_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_test.cpp @@ -1206,8 +1206,7 @@ TEST_F(HybridScanTest, RowGroupPassesMatchesChunkedReader) *footer_buffer, options); auto const all_row_groups = reader->all_row_groups(options); - auto const passes = - reader->construct_row_group_passes(all_row_groups, options, pass_read_limit); + auto const passes = reader->construct_row_group_passes(all_row_groups, pass_read_limit); for (auto const& pass_row_groups : passes) { auto const chunk_byte_ranges = @@ -1219,15 +1218,7 @@ TEST_F(HybridScanTest, RowGroupPassesMatchesChunkedReader) reader->setup_chunking_for_all_columns( 0, pass_read_limit, pass_row_groups, col_data, options, stream, mr); - // Pass planning must not disturb an in-flight chunked materialization. - // `materialize_all_columns_chunk()` does not re-select columns, so a planner call that - // mutated the reader's column selection would silently swap this loop onto a single column. - auto const probe_options = - cudf::io::parquet_reader_options::builder().column_indices({0}).build(); - while (reader->has_next_table_chunk()) { - static_cast( - reader->construct_row_group_passes(pass_row_groups, probe_options, pass_read_limit)); auto chunk = reader->materialize_all_columns_chunk(); hybrid_scan_tables.push_back(std::move(chunk.tbl)); } diff --git a/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp b/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp index d41398610adb..f6f671ae4e7d 100644 --- a/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp +++ b/java/src/main/native/src/HybridScanReaderJniMaterialize.cpp @@ -364,9 +364,8 @@ JNIEXPORT jobjectArray JNICALL Java_ai_rapids_cudf_HybridScanReader_constructRow auto const pass_limit = checked_size_t(env, pass_read_limit, "pass_read_limit"); auto* wrapper = reinterpret_cast(handle); auto holder = make_row_group_span(env, j_row_groups); - auto passes = - wrapper->reader->construct_row_group_passes(holder.span(), wrapper->options, pass_limit); - jclass int_array_cls = env->FindClass("[I"); + auto passes = wrapper->reader->construct_row_group_passes(holder.span(), pass_limit); + jclass int_array_cls = env->FindClass("[I"); if (int_array_cls == nullptr) { return nullptr; } auto outer = env->NewObjectArray(passes.size(), int_array_cls, nullptr); if (outer == nullptr) { return nullptr; } diff --git a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi index 11edfaa0ef0a..f95dc8b054d3 100644 --- a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi +++ b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyi @@ -139,7 +139,6 @@ class HybridScanReader: def construct_row_group_passes( self, row_group_indices: list[int], - options: ParquetReaderOptions, pass_read_limit: int, ) -> list[list[int]]: ... def has_next_table_chunk(self) -> bool: ... diff --git a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx index 2bde59ba088a..664eb489428c 100644 --- a/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx +++ b/python/pylibcudf/pylibcudf/io/experimental/hybrid_scan.pyx @@ -789,7 +789,6 @@ cdef class HybridScanReader: def construct_row_group_passes( self, list row_group_indices, - ParquetReaderOptions options, size_t pass_read_limit, ): """Partition row groups into passes such that the GPU memory required to @@ -797,15 +796,12 @@ cdef class HybridScanReader: Note that ``pass_read_limit`` is a hint, not an absolute limit. i.e. if a row group cannot fit within the limit, it will still constitute a valid - pass. Row group sizes are estimated over only the columns selected for - reading by ``options``. + pass. Parameters ---------- row_group_indices : list[int] Input row group indices - options : ParquetReaderOptions - Parquet reader options, used to determine the columns that will be read pass_read_limit : int Limit on the amount of memory used for reading and decompressing data or 0 if there is no limit. @@ -825,7 +821,6 @@ cdef class HybridScanReader: std_span[const_size_type]( indices_vec.data(), indices_vec.size() ), - options.c_obj, pass_read_limit ) diff --git a/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd b/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd index 4be97f006f94..36201d545de6 100644 --- a/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd +++ b/python/pylibcudf/pylibcudf/libcudf/io/hybrid_scan.pxd @@ -170,7 +170,6 @@ cdef extern from "cudf/io/experimental/hybrid_scan.hpp" \ vector[vector[size_type]] construct_row_group_passes( std_span[const_size_type] row_group_indices, - const parquet_reader_options& options, size_t pass_read_limit, ) except +libcudf_exception_handler diff --git a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py index f0a874ab09ed..74f467f16193 100644 --- a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py +++ b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py @@ -680,21 +680,21 @@ def test_hybrid_scan_construct_row_group_passes( # zero pass read limit => single pass with all row groups pass_read_limit = 0 passes = simple_hybrid_scan_reader.construct_row_group_passes( - all_row_groups, simple_parquet_options, pass_read_limit + all_row_groups, pass_read_limit ) assert passes == [all_row_groups] # small pass read limit => each row group in its own pass pass_read_limit = 1 passes = simple_hybrid_scan_reader.construct_row_group_passes( - all_row_groups, simple_parquet_options, pass_read_limit + all_row_groups, pass_read_limit ) assert passes == [[rg] for rg in all_row_groups] # Passes should flatten to all row groups pass_read_limit = 1024 passes = simple_hybrid_scan_reader.construct_row_group_passes( - all_row_groups, simple_parquet_options, pass_read_limit + all_row_groups, pass_read_limit ) assert [rg for p in passes for rg in p] == all_row_groups assert all(passes) @@ -704,7 +704,7 @@ def test_hybrid_scan_construct_row_group_passes( ValueError, match="Empty input row group indices encountered" ): simple_hybrid_scan_reader.construct_row_group_passes( - [], simple_parquet_options, pass_read_limit + [], pass_read_limit ) From a0e3a147550da54031ddef77cd51915febbd2be4 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Mon, 27 Jul 2026 23:32:26 +0000 Subject: [PATCH 4/9] Minor --- cpp/src/io/parquet/reader_impl_chunking.cu | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_chunking.cu b/cpp/src/io/parquet/reader_impl_chunking.cu index bba5f5038192..9f80a78cd3db 100644 --- a/cpp/src/io/parquet/reader_impl_chunking.cu +++ b/cpp/src/io/parquet/reader_impl_chunking.cu @@ -543,10 +543,7 @@ void reader_impl::compute_input_passes(read_mode mode) ? static_cast(_input_pass_read_limit * input_limit_compression_reserve) : std::numeric_limits::max(); - // Size each row group by the column chunks we are actually going to read. - // `create_global_chunk_info()` has already built one `ColumnChunkDesc` per (row group, selected - // input column), so this is an exact measure of the compressed bytes a pass will hold rather than - // an estimate over the whole row group. + // Compute size information for each row group by the column chunks we are actually going to read auto const row_group_sizes = compute_row_group_size_info(row_groups_info, _file_itm_data.chunks, _input_columns.size()); From f35b314d68f57b055e4b44e453df3740fda82191 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Mon, 27 Jul 2026 23:54:40 +0000 Subject: [PATCH 5/9] Simplify --- .../parquet/experimental/hybrid_scan_impl.cpp | 11 ++-- cpp/src/io/parquet/reader_impl_chunking.cu | 13 +++-- .../io/parquet/reader_impl_chunking_utils.cu | 32 ------------ .../io/parquet/reader_impl_chunking_utils.cuh | 13 ----- cpp/src/io/parquet/reader_impl_helpers.cpp | 52 +++++++++++-------- cpp/src/io/parquet/reader_impl_helpers.hpp | 28 ++++++---- 6 files changed, 60 insertions(+), 89 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 1e3a309bb37c..fc5473b18f11 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -809,14 +809,11 @@ hybrid_scan_reader_impl::construct_row_group_passes( cuda::counting_iterator(row_group_indices.size()), [&](auto const source_index) { for (auto const rg_index : row_group_indices[source_index]) { - auto const& row_group = - _extended_metadata->get_row_group(rg_index, source_index); - auto const [compressed_size, num_rows, max_leaf_values] = - _extended_metadata->get_row_group_properties(row_group); row_group_ids.emplace_back(rg_index, source_index); - row_group_sizes.push_back({.unadjusted_num_rows = num_rows, - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}); + // TODO(mh): Compute the row group size information over the selected columns + // instead + row_group_sizes.push_back(_extended_metadata->get_row_group_size_info( + rg_index, source_index, std::nullopt)); } }); diff --git a/cpp/src/io/parquet/reader_impl_chunking.cu b/cpp/src/io/parquet/reader_impl_chunking.cu index 9f80a78cd3db..6f1c9c05e3cf 100644 --- a/cpp/src/io/parquet/reader_impl_chunking.cu +++ b/cpp/src/io/parquet/reader_impl_chunking.cu @@ -543,9 +543,16 @@ void reader_impl::compute_input_passes(read_mode mode) ? static_cast(_input_pass_read_limit * input_limit_compression_reserve) : std::numeric_limits::max(); - // Compute size information for each row group by the column chunks we are actually going to read - auto const row_group_sizes = - compute_row_group_size_info(row_groups_info, _file_itm_data.chunks, _input_columns.size()); + // Compute size information for each row group by the columns we are actually going to read. + auto row_group_sizes = std::vector{}; + row_group_sizes.reserve(row_groups_info.size()); + std::transform(row_groups_info.cbegin(), + row_groups_info.cend(), + std::back_inserter(row_group_sizes), + [&](auto const& row_group) { + return _metadata->get_row_group_size_info( + row_group.index, row_group.source_index, _input_columns); + }); auto pass_data = compute_row_group_passes(row_group_sizes, comp_read_limit, _file_itm_data.global_skip_rows); diff --git a/cpp/src/io/parquet/reader_impl_chunking_utils.cu b/cpp/src/io/parquet/reader_impl_chunking_utils.cu index 7b7073d05c94..6017d7bfa426 100644 --- a/cpp/src/io/parquet/reader_impl_chunking_utils.cu +++ b/cpp/src/io/parquet/reader_impl_chunking_utils.cu @@ -953,38 +953,6 @@ rmm::device_uvector compute_level_decode_sizes(device_span compute_row_group_size_info( - std::span row_groups_info, - std::span chunks, - size_t num_input_columns) -{ - CUDF_EXPECTS(chunks.size() == row_groups_info.size() * num_input_columns, - "Mismatch between the number of column chunk descriptors and row groups"); - - std::vector row_group_sizes; - row_group_sizes.reserve(row_groups_info.size()); - - std::transform(cuda::counting_iterator(0), - cuda::counting_iterator(row_groups_info.size()), - std::back_inserter(row_group_sizes), - [&](auto const rg_idx) { - auto const& rg_info = row_groups_info[rg_idx]; - auto const chunks_begin = chunks.begin() + (rg_idx * num_input_columns); - auto const chunks_end = chunks_begin + num_input_columns; - size_t compressed_size = 0; - size_t max_leaf_values = 0; - std::for_each(chunks_begin, chunks_end, [&](auto const& chunk) { - compressed_size += chunk.compressed_size; - max_leaf_values = std::max(max_leaf_values, chunk.num_values); - }); - return row_group_size_info{.unadjusted_num_rows = rg_info.unadjusted_num_rows, - .compressed_size = compressed_size, - .max_leaf_values = max_leaf_values}; - }); - - return row_group_sizes; -} - row_group_pass_data compute_row_group_passes(std::span row_group_sizes, std::size_t comp_read_limit, int64_t skip_rows) diff --git a/cpp/src/io/parquet/reader_impl_chunking_utils.cuh b/cpp/src/io/parquet/reader_impl_chunking_utils.cuh index 49cdbfd7e6cc..78026fec76ea 100644 --- a/cpp/src/io/parquet/reader_impl_chunking_utils.cuh +++ b/cpp/src/io/parquet/reader_impl_chunking_utils.cuh @@ -797,17 +797,4 @@ row_group_pass_data compute_row_group_passes(std::span compute_row_group_size_info( - std::span row_groups_info, - std::span chunks, - size_t num_input_columns); - } // namespace cudf::io::parquet::detail diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 4802deabc653..a07b716399f3 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -1209,6 +1209,35 @@ RowGroup const& aggregate_reader_metadata::get_row_group(size_type row_group_ind return per_file_metadata[src_idx].row_groups[row_group_index]; } +row_group_size_info aggregate_reader_metadata::get_row_group_size_info( + size_type row_group_index, + size_type src_idx, + std::optional> input_columns) const +{ + auto const& row_group = get_row_group(row_group_index, src_idx); + + auto size_info = + row_group_size_info{.unadjusted_num_rows = static_cast(row_group.num_rows)}; + + if (input_columns.has_value()) { + for (auto const& column : *input_columns) { + auto const& column_metadata = + get_column_metadata(row_group_index, src_idx, column.schema_idx); + size_info.compressed_size += column_metadata.total_compressed_size; + size_info.max_leaf_values = + std::max(size_info.max_leaf_values, static_cast(column_metadata.num_values)); + } + } else { + for (auto const& column_chunk : row_group.columns) { + size_info.compressed_size += column_chunk.meta_data.total_compressed_size; + size_info.max_leaf_values = + std::max(size_info.max_leaf_values, static_cast(column_chunk.meta_data.num_values)); + } + } + + return size_info; +} + ColumnChunkMetaData const& aggregate_reader_metadata::get_column_metadata(size_type row_group_index, size_type src_idx, int schema_idx) const @@ -1384,29 +1413,6 @@ std::vector aggregate_reader_metadata::get_pandas_index_names() con return names; } -std::tuple aggregate_reader_metadata::get_row_group_properties( - RowGroup const& row_group) const -{ - auto const compressed_size = std::transform_reduce( - row_group.columns.cbegin(), - row_group.columns.cend(), - size_t{0}, - std::plus<>(), - [](auto const& colchunk) { return colchunk.meta_data.total_compressed_size; }); - - size_t const max_leaf_values = - row_group.columns.empty() - ? 0 - : std::max_element(row_group.columns.cbegin(), - row_group.columns.cend(), - [](auto const& a, auto const& b) { - return a.meta_data.num_values < b.meta_data.num_values; - }) - ->meta_data.num_values; - - return {compressed_size, static_cast(row_group.num_rows), max_leaf_values}; -} - std::tuple>, std::vector>, diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index ee53308624a7..52e90917be27 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -16,6 +16,7 @@ #include #include #include +#include #include #include #include @@ -435,6 +436,22 @@ class aggregate_reader_metadata { [[nodiscard]] RowGroup const& get_row_group(size_type row_group_index, size_type src_idx) const; + /** + * @brief Computes row group size information over selected columns + * + * When `input_columns` is specified, computes the compressed size and maximum leaf value count + * over only those columns. Otherwise, over all columns in the row group. + * + * @param row_group_index Index of the row group within its source + * @param src_idx Index of the input source + * @param input_columns Optional selected leaf columns + * @return Row group size information + */ + [[nodiscard]] row_group_size_info get_row_group_size_info( + size_type row_group_index, + size_type src_idx, + std::optional> input_columns) const; + /** * @brief Get Parquet file metadatas * @@ -613,17 +630,6 @@ class aggregate_reader_metadata { */ [[nodiscard]] std::vector get_pandas_index_names() const; - /** - * @brief Computes the compressed size, the number of rows, and the maximum number of leaf values - * in the specified row group over all of its columns - * - * @param row_group Input row group - * - * @return A tuple of row group compressed size, number of rows, and maximum leaf values - */ - [[nodiscard]] std::tuple get_row_group_properties( - RowGroup const& row_group) const; - /** * @brief Filters the row groups using stats and bloom filters based on predicate filter * From cbb013bee2c76657e8d71a2146c8dc877c706036 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 28 Jul 2026 00:32:00 +0000 Subject: [PATCH 6/9] Improve test --- cpp/tests/io/parquet_chunked_reader_test.cu | 131 +++++++------------- 1 file changed, 44 insertions(+), 87 deletions(-) diff --git a/cpp/tests/io/parquet_chunked_reader_test.cu b/cpp/tests/io/parquet_chunked_reader_test.cu index 9c4e762e1033..4c9e7d056ed2 100644 --- a/cpp/tests/io/parquet_chunked_reader_test.cu +++ b/cpp/tests/io/parquet_chunked_reader_test.cu @@ -103,14 +103,8 @@ auto write_file(std::vector>& input_columns, return std::pair{std::move(input_table), std::move(filepath)}; } -auto chunked_read(std::vector const& filepaths, - std::size_t output_limit, - std::size_t input_limit = 0) +auto chunked_read(cudf::io::chunked_parquet_reader const& reader) { - auto const read_opts = - cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepaths}).build(); - auto reader = cudf::io::chunked_parquet_reader(output_limit, input_limit, read_opts); - auto num_chunks = 0; auto out_tables = std::vector>{}; @@ -133,6 +127,23 @@ auto chunked_read(std::vector const& filepaths, return std::pair(cudf::concatenate(out_tviews), num_chunks); } +auto chunked_read(cudf::io::parquet_reader_options const& read_opts, + std::size_t output_limit, + std::size_t input_limit = 0) +{ + auto reader = cudf::io::chunked_parquet_reader(output_limit, input_limit, read_opts); + return chunked_read(reader); +} + +auto chunked_read(std::vector const& filepaths, + std::size_t output_limit, + std::size_t input_limit = 0) +{ + auto const read_opts = + cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepaths}).build(); + return chunked_read(read_opts, output_limit, input_limit); +} + auto chunked_read(std::string const& filepath, std::size_t output_limit, std::size_t input_limit = 0) @@ -149,27 +160,7 @@ auto chunked_read(std::vector>&& sources, auto const read_opts = cudf::io::parquet_reader_options::builder().build(); auto reader = cudf::io::chunked_parquet_reader( output_limit, input_limit, std::move(sources), std::move(metadatas), read_opts); - - auto num_chunks = 0; - auto out_tables = std::vector>{}; - - do { - auto chunk = reader.read_chunk(); - // If the input file is empty, the first call to `read_chunk` will return an empty table. - // Thus, we only check for non-empty output table from the second call. - if (num_chunks > 0) { - CUDF_EXPECTS(chunk.tbl->num_rows() != 0, "Number of rows in the new chunk is zero."); - } - ++num_chunks; - out_tables.emplace_back(std::move(chunk.tbl)); - } while (reader.has_next()); - - auto out_tviews = std::vector{}; - for (auto const& tbl : out_tables) { - out_tviews.emplace_back(tbl->view()); - } - - return std::pair(cudf::concatenate(out_tviews), num_chunks); + return chunked_read(reader); } auto const read_table_and_nrows_per_source(cudf::io::chunked_parquet_reader const& reader) @@ -1338,85 +1329,51 @@ TEST_F(ParquetChunkedReaderInputLimitConstrainedTest, MixedColumns) struct ParquetChunkedReaderInputLimitTest : public cudf::test::BaseFixture {}; -namespace { -// Reads a file through the chunked reader, optionally projecting a subset of columns. -auto chunked_read_projected(std::string const& filepath, - std::vector const& columns, - std::size_t output_limit, - std::size_t input_limit) +TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsReducePasses) { - auto builder = cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}); - if (not columns.empty()) { builder.column_names(columns); } - auto const read_opts = builder.build(); - auto reader = cudf::io::chunked_parquet_reader(output_limit, input_limit, read_opts); + constexpr int num_columns = 16; + constexpr int num_rows = 2'000; + constexpr int rows_per_row_group = num_rows / 2; - auto num_chunks = 0; - auto out_tables = std::vector>{}; - do { - auto chunk = reader.read_chunk(); - ++num_chunks; - out_tables.emplace_back(std::move(chunk.tbl)); - } while (reader.has_next()); - - auto out_tviews = std::vector{}; - for (auto const& tbl : out_tables) { - out_tviews.emplace_back(tbl->view()); - } - return std::pair(cudf::concatenate(out_tviews), num_chunks); -} -} // namespace - -// Pass construction must size each row group by the columns actually being read. Projecting a -// single column out of a wide file should therefore need far fewer passes - and so produce far -// fewer output chunks - than reading every column under the same `pass_read_limit`. -TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsShrinkPasses) -{ - constexpr int num_columns = 16; - constexpr int num_rows = 125'000; - - auto const filepath = temp_env->get_temp_filepath("projected_columns_shrink_passes.parquet"); + auto const filepath = temp_env->get_temp_filepath("ProjectedColumnPasses.parquet"); auto const iter = cuda::counting_iterator{0}; auto columns = std::vector>{}; auto column_names = std::vector{}; for (int c = 0; c < num_columns; c++) { - auto col = cudf::test::fixed_width_column_wrapper(iter, iter + num_rows); - columns.emplace_back(col.release()); - column_names.emplace_back("col_" + std::to_string(c)); + columns.emplace_back( + cudf::test::fixed_width_column_wrapper(iter, iter + num_rows).release()); } - auto const input_table = cudf::table(std::move(columns)); - auto metadata = cudf::io::table_input_metadata{input_table.view()}; - for (int c = 0; c < num_columns; c++) { - metadata.column_metadata[c].set_name(column_names[c]); - } - - // Many small row groups so that the pass builder has boundaries to choose between, and no - // dictionary encoding so that column chunk sizes are predictable. + auto const input_table = cudf::table(std::move(columns)); cudf::io::write_parquet( cudf::io::parquet_writer_options::builder(cudf::io::sink_info{filepath}, input_table.view()) - .metadata(std::move(metadata)) + .metadata(cudf::io::table_input_metadata{input_table.view()}) .compression(cudf::io::compression_type::NONE) .dictionary_policy(cudf::io::dictionary_policy::NEVER) - .row_group_size_rows(5'000) + .max_page_fragment_size(rows_per_row_group) + .row_group_size_rows(rows_per_row_group) .build()); - // 25 row groups, each column chunk 5'000 * sizeof(int32_t) == 20'000 bytes, so a row group is - // ~320'000 bytes over all 16 columns. The pass budget is - // `pass_read_limit * input_limit_compression_reserve` == ~400'000 bytes: one row group per pass - // for the full read, but twenty row groups per pass for a single column. - constexpr std::size_t pass_read_limit = 1'333'334; + // Each row group is roughly 64 KB = sizeof(int32) * 1000 * 16. Use 100KB as limit so we pick one + // row group per pass when all columns are read but more than one for single column read. + constexpr std::size_t pass_read_limit = 100'000; - auto const [all_cols, all_cols_chunks] = chunked_read_projected(filepath, {}, 0, pass_read_limit); - auto const [one_col, one_col_chunks] = - chunked_read_projected(filepath, {column_names[0]}, 0, pass_read_limit); + auto const all_columns_options = + cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}).build(); + auto const [all_cols, all_cols_chunks] = chunked_read(all_columns_options, 0, pass_read_limit); - // The projected read must not be split more finely than the full read. - EXPECT_LT(one_col_chunks, all_cols_chunks); + auto const one_column_options = + cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}) + .column_indices({0}) + .build(); + auto const [one_col, one_col_chunks] = chunked_read(one_column_options, 0, pass_read_limit); - // ...and it must still return the right data. CUDF_TEST_EXPECT_TABLES_EQUAL(all_cols->view(), input_table.view()); CUDF_TEST_EXPECT_TABLES_EQUAL(one_col->view(), cudf::table_view{{input_table.view().column(0)}}); + + // Projected read must yield fewer chunks than full read + EXPECT_LT(one_col_chunks, all_cols_chunks); } namespace { From 993e2a10cc02cf783c8fefaea919c13f6650c96c Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 28 Jul 2026 00:36:14 +0000 Subject: [PATCH 7/9] Improve tests --- cpp/tests/io/parquet_chunked_reader_test.cu | 38 +++++++++++---------- 1 file changed, 20 insertions(+), 18 deletions(-) diff --git a/cpp/tests/io/parquet_chunked_reader_test.cu b/cpp/tests/io/parquet_chunked_reader_test.cu index 4c9e7d056ed2..4e259ab719d1 100644 --- a/cpp/tests/io/parquet_chunked_reader_test.cu +++ b/cpp/tests/io/parquet_chunked_reader_test.cu @@ -1334,18 +1334,21 @@ TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsReducePasses) constexpr int num_columns = 16; constexpr int num_rows = 2'000; constexpr int rows_per_row_group = num_rows / 2; - - auto const filepath = temp_env->get_temp_filepath("ProjectedColumnPasses.parquet"); - - auto const iter = cuda::counting_iterator{0}; - auto columns = std::vector>{}; - auto column_names = std::vector{}; - for (int c = 0; c < num_columns; c++) { - columns.emplace_back( - cudf::test::fixed_width_column_wrapper(iter, iter + num_rows).release()); - } - + auto const filepath = temp_env->get_temp_filepath("ProjectedColumnPasses.parquet"); + + // Input table + auto columns = std::vector>{}; + std::transform(cuda::counting_iterator{0}, + cuda::counting_iterator{num_columns}, + std::back_inserter(columns), + [](int c) { + return cudf::test::fixed_width_column_wrapper( + cuda::counting_iterator{0}, cuda::counting_iterator{num_rows}) + .release(); + }); auto const input_table = cudf::table(std::move(columns)); + + // Write table to parquet cudf::io::write_parquet( cudf::io::parquet_writer_options::builder(cudf::io::sink_info{filepath}, input_table.view()) .metadata(cudf::io::table_input_metadata{input_table.view()}) @@ -1359,15 +1362,14 @@ TEST_F(ParquetChunkedReaderInputLimitTest, ProjectedColumnsReducePasses) // row group per pass when all columns are read but more than one for single column read. constexpr std::size_t pass_read_limit = 100'000; - auto const all_columns_options = + auto const all_columns = cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}).build(); - auto const [all_cols, all_cols_chunks] = chunked_read(all_columns_options, 0, pass_read_limit); + auto const [all_cols, all_cols_chunks] = chunked_read(all_columns, 0, pass_read_limit); - auto const one_column_options = - cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}) - .column_indices({0}) - .build(); - auto const [one_col, one_col_chunks] = chunked_read(one_column_options, 0, pass_read_limit); + auto const one_column = cudf::io::parquet_reader_options::builder(cudf::io::source_info{filepath}) + .column_indices({0}) + .build(); + auto const [one_col, one_col_chunks] = chunked_read(one_column, 0, pass_read_limit); CUDF_TEST_EXPECT_TABLES_EQUAL(all_cols->view(), input_table.view()); CUDF_TEST_EXPECT_TABLES_EQUAL(one_col->view(), cudf::table_view{{input_table.view().column(0)}}); From 65275ae543b82eaea751a48fb263ebbf14453cc5 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Mon, 3 Aug 2026 19:25:03 +0000 Subject: [PATCH 8/9] style fix --- cpp/src/io/parquet/reader_impl_helpers.hpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index f7b228a0e6ce..ed67ca05cfc3 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -457,7 +457,7 @@ class aggregate_reader_metadata { size_type row_group_index, size_type src_idx, std::optional> input_columns) const; - + /** * @brief Check if all row groups have an offset index * From 709a5fd16ad20b119fc53cd7cf0a1cf173375bb7 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 4 Aug 2026 02:35:59 +0000 Subject: [PATCH 9/9] Minor overflow check --- cpp/src/io/parquet/reader_impl_helpers.cpp | 22 ++++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 4391b38461fc..4d1ecc404376 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -1243,22 +1243,36 @@ row_group_size_info aggregate_reader_metadata::get_row_group_size_info( { auto const& row_group = get_row_group(row_group_index, src_idx); + CUDF_EXPECTS(row_group.num_rows >= 0 && std::in_range(row_group.num_rows), + "Row group has an invalid number of rows", + std::invalid_argument); auto size_info = row_group_size_info{.unadjusted_num_rows = static_cast(row_group.num_rows)}; + // Helper function to overflow-safeadd compressed sizes + auto const add_compressed_size = [&](auto&& current_size, auto&& colchunk_compressed_size) { + auto const sum = cuda::add_overflow(current_size, colchunk_compressed_size); + CUDF_EXPECTS(not sum.overflow, + "Row group compressed size exceeds the supported range", + std::overflow_error); + return sum.value; + }; + if (input_columns.has_value()) { for (auto const& column : *input_columns) { auto const& column_metadata = get_column_metadata(row_group_index, src_idx, column.schema_idx); - size_info.compressed_size += column_metadata.total_compressed_size; + size_info.compressed_size = + add_compressed_size(size_info.compressed_size, column_metadata.total_compressed_size); size_info.max_leaf_values = - std::max(size_info.max_leaf_values, static_cast(column_metadata.num_values)); + std::max(size_info.max_leaf_values, column_metadata.num_values); } } else { for (auto const& column_chunk : row_group.columns) { - size_info.compressed_size += column_chunk.meta_data.total_compressed_size; + size_info.compressed_size = add_compressed_size(size_info.compressed_size, + column_chunk.meta_data.total_compressed_size); size_info.max_leaf_values = - std::max(size_info.max_leaf_values, static_cast(column_chunk.meta_data.num_values)); + std::max(size_info.max_leaf_values, column_chunk.meta_data.num_values); } }