From b8749ad658b65f17f026d8df4dedaaa104bc013e Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Wed, 19 Aug 2026 21:59:51 +0000 Subject: [PATCH 1/3] Add fallback page pruning method --- .../hybrid_scan_io/hybrid_scan_composer.cpp | 24 ++-- .../experimental/hybrid_scan_chunking.cu | 10 +- .../parquet/experimental/hybrid_scan_impl.cpp | 74 +++++++++- .../parquet/experimental/hybrid_scan_impl.hpp | 6 + .../parquet/experimental/page_index_filter.cu | 134 ++++++++++-------- .../experimental/page_index_filter_utils.hpp | 22 ++- .../experimental/hybrid_scan_filters_test.cpp | 2 +- 7 files changed, 188 insertions(+), 84 deletions(-) diff --git a/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp b/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp index a26839d9ffd6..5aa2017da3b5 100644 --- a/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp +++ b/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp @@ -426,22 +426,6 @@ std::unique_ptr hybrid_scan( } } -// Specialization for two-step read without page index -template - requires(not single_step_read and not use_page_index) -std::unique_ptr inline hybrid_scan( - io_source const& io_source, - std::optional filter_expression, - std::unordered_set const& filters, - bool verbose, - rmm::cuda_stream_view stream, - rmm::device_async_resource_ref mr) -{ - static_assert(single_step_read or use_page_index, - "Hybrid scan requires parquet page index for two-step parquet read"); - return nullptr; -} - // Instantiations for hybrid_scan template template std::unique_ptr hybrid_scan( @@ -460,6 +444,14 @@ template std::unique_ptr hybrid_scan( rmm::cuda_stream_view, rmm::device_async_resource_ref); +template std::unique_ptr hybrid_scan( + io_source const&, + std::optional, + std::unordered_set const&, + bool, + rmm::cuda_stream_view, + rmm::device_async_resource_ref); + template std::unique_ptr hybrid_scan( io_source const&, std::optional, diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_chunking.cu b/cpp/src/io/parquet/experimental/hybrid_scan_chunking.cu index 99728cc00ae1..81952a8b9f14 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_chunking.cu +++ b/cpp/src/io/parquet/experimental/hybrid_scan_chunking.cu @@ -40,8 +40,16 @@ void hybrid_scan_reader_impl::handle_chunking( // setup the next pass setup_next_pass(column_chunk_data); + // Compute the data page mask from decoded page headers if needed + auto const data_page_mask_pghdr = [&]() { + if (not _has_offset_index and not _row_mask.is_empty()) { + return compute_data_page_mask_with_page_headers(); + } + return thrust::host_vector(data_page_mask.begin(), data_page_mask.end()); + }(); + // Must be called as soon as we create the pass - set_pass_page_mask(data_page_mask); + set_pass_page_mask(data_page_mask_pghdr.empty() ? data_page_mask : data_page_mask_pghdr); } auto& pass = *_pass_itm_data; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index c8259909e106..418db8b0c018 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -8,6 +8,7 @@ #include "cudf/io/text/byte_range_info.hpp" #include "hybrid_scan_helpers.hpp" #include "io/parquet/reader_impl_chunking_utils.cuh" +#include "page_index_filter_utils.hpp" #include #include @@ -545,6 +546,7 @@ table_with_metadata hybrid_scan_reader_impl::materialize_filter_columns( auto data_page_mask = thrust::host_vector{}; if (mask_data_pages == use_data_page_mask::YES) { + _row_mask = row_mask; data_page_mask = _extended_metadata->compute_data_page_mask( row_mask, row_group_indices, _input_columns, _row_mask_offset, stream); } @@ -583,6 +585,7 @@ table_with_metadata hybrid_scan_reader_impl::materialize_payload_columns( auto data_page_mask = thrust::host_vector{}; if (not row_mask.is_empty() and mask_data_pages == use_data_page_mask::YES) { + _row_mask = row_mask; data_page_mask = _extended_metadata->compute_data_page_mask( row_mask, row_group_indices, _input_columns, _row_mask_offset, stream); } @@ -654,6 +657,7 @@ void hybrid_scan_reader_impl::setup_chunking_for_filter_columns( auto data_page_mask = thrust::host_vector{}; if (mask_data_pages == use_data_page_mask::YES) { + _row_mask = row_mask; data_page_mask = _extended_metadata->compute_data_page_mask( row_mask, row_group_indices, _input_columns, _row_mask_offset, stream); } @@ -714,6 +718,7 @@ void hybrid_scan_reader_impl::setup_chunking_for_payload_columns( auto data_page_mask = thrust::host_vector{}; if (not row_mask.is_empty() and mask_data_pages == use_data_page_mask::YES) { + _row_mask = row_mask; data_page_mask = _extended_metadata->compute_data_page_mask( row_mask, row_group_indices, _input_columns, _row_mask_offset, stream); } @@ -873,7 +878,6 @@ bool hybrid_scan_reader_impl::has_next_table_chunk() void hybrid_scan_reader_impl::reset_internal_state() { - _row_mask_offset = 0; _file_itm_data = file_intermediate_data{}; _file_preprocessed = false; _has_offset_index = false; @@ -896,6 +900,10 @@ void hybrid_scan_reader_impl::reset_internal_state() _output_chunk_read_limit = 0; _strings_to_categorical = false; _reader_column_schema.reset(); + + _row_mask = column_view{}; + _row_mask_offset = 0; + _expr_conv = named_to_reference_converter{}; _mr = cudf::get_current_device_resource_ref(); } @@ -1239,4 +1247,68 @@ void hybrid_scan_reader_impl::set_pass_page_mask(std::span data_page "Encountered mismatch in number of pass pages and page mask size"); } +thrust::host_vector hybrid_scan_reader_impl::compute_data_page_mask_with_page_headers() +{ + auto& pass = *_pass_itm_data; + pass.pages.device_to_host_async(_stream); + _stream.sync(); + + std::vector page_row_offsets; + page_row_offsets.reserve(pass.pages.size() * 2); + + // Maps each data page to its flat-page range; -1 keeps nested pages enabled. + std::vector row_range_map; + row_range_map.reserve(pass.pages.size()); + + cudf::size_type previous_chunk_idx = -1; + auto max_page_size = cudf::size_type{0}; + + for (auto const& page : pass.pages) { + // Ignore dictionary pages altogether + if (page.flags & parquet::detail::PAGEINFO_FLAGS_DICTIONARY) { continue; } + + auto const& chunk = pass.chunks[page.chunk_idx]; + + // Don't prune list column pages as rows may span page boundaries when offset index isn't + // present. + if (chunk.max_level[parquet::detail::level_type::REPETITION] > 0) { + row_range_map.push_back(-1); + continue; + } + + auto const page_start = chunk.start_row + page.chunk_row; + auto const page_end = page_start + page.num_rows; + max_page_size = std::max(max_page_size, page_end - page_start); + + // Starting a new column chunk. Push page start row + if (previous_chunk_idx == -1 or page.chunk_idx != previous_chunk_idx) { + page_row_offsets.push_back(page_start); + previous_chunk_idx = page.chunk_idx; + } + + // Push row range index and page end row + row_range_map.push_back(page_row_offsets.size() - 1); + page_row_offsets.push_back(page_end); + } + + auto data_page_mask = thrust::host_vector{}; + + // Compute the row range mask + auto const row_range_mask = compute_row_range_selection_mask( + _row_mask, _row_mask_offset, pass.num_rows, page_row_offsets, max_page_size, _stream); + + if (row_range_mask.empty()) { return data_page_mask; } + + CUDF_EXPECTS(row_range_mask.size() == page_row_offsets.size() - 1, + "Encountered invalid row range mask size"); + + data_page_mask.reserve(row_range_map.size()); + + // Scatter row range results while retaining list column pages. + for (auto const range_idx : row_range_map) { + data_page_mask.push_back(range_idx < 0 ? true : row_range_mask[range_idx]); + } + return data_page_mask; +} + } // namespace cudf::io::parquet::experimental::detail diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index e0e7d1160352..91208732806c 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -357,6 +357,11 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { */ void set_pass_page_mask(std::span data_page_mask); + /** + * @brief Compute a data page mask from the decoded page headers. + */ + [[nodiscard]] thrust::host_vector compute_data_page_mask_with_page_headers(); + /** * @brief Select the columns to be read based on the read mode * @@ -567,6 +572,7 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { std::optional> _filter_columns_names; + cudf::column_view _row_mask{}; cudf::size_type _row_mask_offset{0}; bool _output_chunk_produced{false}; diff --git a/cpp/src/io/parquet/experimental/page_index_filter.cu b/cpp/src/io/parquet/experimental/page_index_filter.cu index 61312327d8ac..bf5c8c095ee1 100644 --- a/cpp/src/io/parquet/experimental/page_index_filter.cu +++ b/cpp/src/io/parquet/experimental/page_index_filter.cu @@ -976,6 +976,68 @@ std::unique_ptr aggregate_reader_metadata::build_row_mask_with_pag page_stats_table, stats_expr.get_stats_expr().get(), stream, mr); } +thrust::host_vector compute_row_range_selection_mask( + cudf::column_view const& row_mask, + cudf::size_type row_mask_offset, + cudf::size_type total_rows, + std::span page_row_offsets, + cudf::size_type max_page_size, + cuda::stream_ref stream) +{ + // Need at least two offsets (or one range) to search the Fenwick tree + if (page_row_offsets.size() < 2) return thrust::host_vector{}; + + // Return early if all rows are needed. + if (cudf::detail::all_of(row_mask.begin() + row_mask_offset, + row_mask.begin() + row_mask_offset + total_rows, + cuda::std::identity{}, + stream)) { + return thrust::host_vector{}; + } + + auto const mr = cudf::get_current_device_resource_ref(); + auto const tree_level_offsets = compute_fenwick_tree_level_offsets(total_rows, max_page_size); + auto const num_levels = static_cast(tree_level_offsets.size()); + auto tree_levels_data = rmm::device_uvector(tree_level_offsets.back(), stream, mr); + auto host_tree_level_ptrs = cudf::detail::make_pinned_vector_async(num_levels, stream); + host_tree_level_ptrs[0] = const_cast(row_mask.begin()) + row_mask_offset; + std::for_each(cuda::counting_iterator{1}, + cuda::counting_iterator{num_levels}, + [&](auto const level_idx) { + host_tree_level_ptrs[level_idx] = + tree_levels_data.data() + tree_level_offsets[level_idx - 1]; + }); + auto tree_level_ptrs = cudf::detail::make_device_uvector_async(host_tree_level_ptrs, stream, mr); + + auto prev_level_size = total_rows; + std::for_each( + cuda::counting_iterator{0}, + cuda::counting_iterator{num_levels - 1}, + [&](auto const prev_level) { + auto const current_level_size = cudf::util::div_rounding_up_safe(prev_level_size, 2); + thrust::for_each(rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()), + cuda::counting_iterator{0}, + cuda::counting_iterator{current_level_size}, + build_fenwick_tree_level_functor{ + tree_level_ptrs.data(), prev_level, prev_level_size, current_level_size}); + prev_level_size = current_level_size; + }); + + auto const num_ranges = static_cast(page_row_offsets.size() - 1); + auto device_results = rmm::device_uvector(num_ranges, stream, mr); + auto pinned_page_offsets = cudf::detail::make_pinned_vector(page_row_offsets, stream); + auto page_offsets = cudf::detail::make_device_uvector_async(pinned_page_offsets, stream, mr); + thrust::transform( + rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()), + cuda::counting_iterator{0}, + cuda::counting_iterator{num_ranges}, + device_results.begin(), + search_fenwick_tree_functor{tree_level_ptrs.data(), page_offsets.data(), num_ranges}); + auto results = cudf::detail::make_pinned_vector_async(device_results, stream); + stream.sync(); + return thrust::host_vector(results.begin(), results.end()); +} + template thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( ColumnView const& row_mask, @@ -1113,78 +1175,24 @@ thrust::host_vector aggregate_reader_metadata::compute_data_page_mask( "Row mask must not contain nulls for payload columns"); } - auto const mr = cudf::get_current_device_resource_ref(); - - // Compute fenwick tree level offsets and total size (level 1 and higher) - auto const tree_level_offsets = compute_fenwick_tree_level_offsets(total_rows, max_page_size); - auto const num_levels = static_cast(tree_level_offsets.size()); - // Buffer to store Fenwick tree levels (level 1 and higher) data - auto tree_levels_data = rmm::device_uvector(tree_level_offsets.back(), stream, mr); - - // Pointers to each Fenwick tree level data - auto host_tree_level_ptrs = cudf::detail::make_pinned_vector_async(num_levels, stream); - // Zeroth level is just the row mask itself - host_tree_level_ptrs[0] = const_cast(row_mask.template begin()) + row_mask_offset; - std::for_each(cuda::counting_iterator{1}, - cuda::counting_iterator{num_levels}, - [&](auto const level_idx) { - host_tree_level_ptrs[level_idx] = - tree_levels_data.data() + tree_level_offsets[level_idx - 1]; - }); - - auto fenwick_tree_level_ptrs = - cudf::detail::make_device_uvector_async(host_tree_level_ptrs, stream, mr); + auto data_page_mask = thrust::host_vector{}; - // Build Fenwick tree levels (zeroth level is just the row mask itself) - auto prev_level_size = static_cast(total_rows); - std::for_each( - cuda::counting_iterator{0}, - cuda::counting_iterator{num_levels - 1}, - [&](auto const prev_level) { - auto const current_level_size = cudf::util::div_rounding_up_safe(prev_level_size, 2); - thrust::for_each( - rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()), - cuda::counting_iterator{0}, - cuda::counting_iterator{current_level_size}, - build_fenwick_tree_level_functor{ - fenwick_tree_level_ptrs.data(), prev_level, prev_level_size, current_level_size}); - prev_level_size = current_level_size; - }); - - // Search the Fenwick tree to see if there's a surviving row in each page's row range - auto const num_ranges = static_cast(page_row_offsets.size() - 1); - rmm::device_uvector device_data_page_mask(num_ranges, stream, mr); - // Use a pinned bounce buffer to avoid pageable h2d copy - auto pinned_page_offsets = cudf::detail::make_pinned_vector( - cudf::host_span{page_row_offsets}, stream); - auto page_offsets = cudf::detail::make_device_uvector_async(pinned_page_offsets, stream, mr); - thrust::transform( - rmm::exec_policy_nosync(stream, cudf::get_current_device_resource_ref()), - cuda::counting_iterator{0}, - cuda::counting_iterator{num_ranges}, - device_data_page_mask.begin(), - search_fenwick_tree_functor{fenwick_tree_level_ptrs.data(), page_offsets.data(), num_ranges}); - - // Copy over search results to host - auto host_results = cudf::detail::make_pinned_vector_async(device_data_page_mask, stream); - auto const total_pages = pinned_page_offsets.size() - num_columns; - auto data_page_mask = thrust::host_vector(total_pages); - auto host_results_iter = host_results.begin(); - stream.sync(); + auto const row_range_mask = compute_row_range_selection_mask( + row_mask, row_mask_offset, total_rows, page_row_offsets, max_page_size, stream); + if (row_range_mask.empty()) { return data_page_mask; } + data_page_mask.reserve(page_row_offsets.size() - num_columns); // Discard results for invalid ranges. i.e. ranges starting at the last page of a column and // ending at the first page of the next column - auto num_pages_inserted = 0; std::for_each(cuda::counting_iterator{0}, cuda::counting_iterator{num_columns}, [&](auto col_idx) { auto const col_num_pages = col_page_offsets[col_idx + 1] - col_page_offsets[col_idx] - 1; - data_page_mask.insert(data_page_mask.begin() + num_pages_inserted, - host_results_iter, - host_results_iter + col_num_pages); - host_results_iter += col_num_pages + 1; - num_pages_inserted += col_num_pages; + auto const first_page_range = col_page_offsets[col_idx]; + data_page_mask.insert(data_page_mask.end(), + row_range_mask.begin() + first_page_range, + row_range_mask.begin() + first_page_range + col_num_pages); }); return data_page_mask; } diff --git a/cpp/src/io/parquet/experimental/page_index_filter_utils.hpp b/cpp/src/io/parquet/experimental/page_index_filter_utils.hpp index 264510ff48fd..b1166d9b574f 100644 --- a/cpp/src/io/parquet/experimental/page_index_filter_utils.hpp +++ b/cpp/src/io/parquet/experimental/page_index_filter_utils.hpp @@ -47,8 +47,7 @@ compute_page_row_offsets_and_colchunk_page_offsets( * @param per_file_metadata Span of parquet footer metadata * @param row_group_indices Span of input row group indices * @param schema_idx Column's schema index - * @return A pair of page row offsets and the size of the largest page in this - * column + * @return A pair of page row offsets and the size of the largest page in this column */ [[nodiscard]] std::pair, size_type> compute_page_row_offsets( cudf::host_span per_file_metadata, @@ -81,4 +80,23 @@ compute_page_row_offsets_and_colchunk_page_offsets( [[nodiscard]] std::vector compute_fenwick_tree_level_offsets( cudf::size_type level0_size, cudf::size_type max_page_size); +/** + * @brief Computes a mask indicating which row ranges contain at least one selected row. + * + * @param row_mask Boolean column indicating selected rows + * @param row_mask_offset Offset into the row mask for the current pass + * @param total_rows Number of rows in the current pass + * @param page_row_offsets Page row offsets defining the row ranges + * @param max_page_size Size of the largest page row range + * @param stream CUDA stream used for device memory operations and kernel launches + * @return Boolean vector with one entry for each consecutive row range + */ +[[nodiscard]] thrust::host_vector compute_row_range_selection_mask( + cudf::column_view const& row_mask, + cudf::size_type row_mask_offset, + cudf::size_type total_rows, + std::span page_row_offsets, + cudf::size_type max_page_size, + cuda::stream_ref stream); + } // namespace cudf::io::parquet::experimental::detail diff --git a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp index eee71f349ec7..5a2f236d7580 100644 --- a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp @@ -906,7 +906,7 @@ TEST_F(HybridScanFiltersTest, OffsetIndexOnlyDataPageMask) auto const expected = cudf::apply_boolean_mask(written_table->view(), row_mask_view, stream, mr); CUDF_TEST_EXPECT_TABLES_EQUIVALENT(expected->view(), result.tbl->view()); - // Without offset index, data-page pruning falls back to decoding all pages. + // Without an offset index, data-page pruning derives page ranges from decoded page headers. for (auto& row_group : metadata.row_groups) { for (auto& column : row_group.columns) { column.offset_index.reset(); From b109b8616db1cab2f8bc12231a34870f74d2a7db Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Wed, 19 Aug 2026 22:28:22 +0000 Subject: [PATCH 2/3] style --- cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp b/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp index 5aa2017da3b5..c9fa43bea8eb 100644 --- a/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp +++ b/cpp/examples/hybrid_scan_io/hybrid_scan_composer.cpp @@ -1,5 +1,5 @@ /* - * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION. + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. * SPDX-License-Identifier: Apache-2.0 */ From 3903c10d3c540fc33ac8ce4c7562c2efb495bc0b Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Wed, 19 Aug 2026 22:29:34 +0000 Subject: [PATCH 3/3] Java and Python changes --- .../java/ai/rapids/cudf/HybridScanReader.java | 18 +-- .../ai/rapids/cudf/HybridScanReaderTest.java | 133 +++++++++++++++--- .../tests/io/test_experimental_hybrid_scan.py | 67 +++++++++ 3 files changed, 191 insertions(+), 27 deletions(-) diff --git a/java/src/main/java/ai/rapids/cudf/HybridScanReader.java b/java/src/main/java/ai/rapids/cudf/HybridScanReader.java index b783aa86cff4..2165abd9a085 100644 --- a/java/src/main/java/ai/rapids/cudf/HybridScanReader.java +++ b/java/src/main/java/ai/rapids/cudf/HybridScanReader.java @@ -39,10 +39,14 @@ * chunked reader pipeline. * *

The filter and payload materialization paths accept a boolean that toggles - * page-level pruning: skips decode of pages the filter (or row mask) proves empty, in - * exchange for a per-page stats scan and a carried row-mask column. Enable when the - * workload prunes many pages; requires prior {@link #setupPageIndex(HostMemoryBuffer)} to - * prune filter column pages using page-level statistics. + * page-level pruning: skips decode of pages the filter (or row mask) proves empty. Enable when the + * workload prunes many pages. Pruning requirements for the two column materializations differ: + *

    + *
  • Filter columns seed the row mask from page-index statistics, so + * {@link #setupPageIndex(HostMemoryBuffer)} must have been called first.
  • + *
  • Payload columns only need page row boundaries. These come from the + * {@code OffsetIndex} when page index has been set up, otherwise from the decoded page headers as fallback.
  • + *
* *

The reader is created with no filter expression installed. Filter-related APIs * behave as though nothing has been filtered out unless a filter is first supplied via @@ -382,8 +386,7 @@ public FilterMaterializationResult materializeFilterColumns(int[] rowGroupIndice * returned by {@link #payloadColumnChunksByteRanges(int[])} * @param rowMask row mask (read-only) * @param usePageLevelPruning enable the data page mask to skip decode of pages the row - * mask proves empty; requires prior - * {@link #setupPageIndex(HostMemoryBuffer)} to avoid fall back path + * mask proves empty * @return the materialized payload column table */ public Table materializePayloadColumns(int[] rowGroupIndices, @@ -529,8 +532,7 @@ public ColumnVector takeFilterRowMask() { * @param rowGroupIndices row groups to read * @param rowMask row mask (read-only) * @param usePageLevelPruning enable the data page mask to skip decode of pages the row - * mask proves empty; requires prior - * {@link #setupPageIndex(HostMemoryBuffer)} to avoid fall back path + * mask proves empty * @param columnChunkData device buffers holding the payload column chunks, in the order * returned by {@link #payloadColumnChunksByteRanges(int[])} */ diff --git a/java/src/test/java/ai/rapids/cudf/HybridScanReaderTest.java b/java/src/test/java/ai/rapids/cudf/HybridScanReaderTest.java index 5eed87e1e0e0..463323e380fe 100644 --- a/java/src/test/java/ai/rapids/cudf/HybridScanReaderTest.java +++ b/java/src/test/java/ai/rapids/cudf/HybridScanReaderTest.java @@ -593,6 +593,42 @@ void testMaterializePayloadColumnsExactRowCount(@TempDir Path tmp) throws IOExce } } + /** + * Verifies materializePayloadColumns() prunes pages from page-header row counts when the + * file has no page index: zip_code > 100,000 keeps only the last of row group 1's three + * pages, so the payload must be exactly the 19,999 rows with ids 100,001–119,999. + */ + @Test + void testMaterializePayloadColumnsPagePruningWithoutPageIndex(@TempDir Path tmp) + throws IOException { + try (OpenReader open = + OpenReader.multiPage(tmp).withFilter("zip_code", BinaryOperator.GREATER, 100000)) { + HybridScanReader reader = open.reader; + assertEquals(0L, reader.pageIndexByteRange().size(), + "Fixture must have no page index so the header-derived fallback is exercised"); + int[] survived = reader.filterRowGroupsWithStats(reader.allRowGroups()); + assertArrayEquals(new int[]{1}, survived, + "Group 0 (zip_code 0–59,999) cannot satisfy zip_code > 100,000"); + DeviceMemoryBuffer[] filterCols = copyRangesToDevice( + open.file, reader.filterColumnChunksByteRanges(survived)); + DeviceMemoryBuffer[] payloadCols = copyRangesToDevice( + open.file, reader.payloadColumnChunksByteRanges(survived)); + try (HybridScanReader.FilterMaterializationResult fr = + reader.materializeFilterColumns(survived, filterCols, false); + Table payload = reader.materializePayloadColumns(survived, payloadCols, + fr.rowMask(), true); + ColumnVector expectedIds = ColumnVector.fromInts( + IntStream.rangeClosed(100001, 119999).toArray())) { + assertEquals(2, payload.getNumberOfColumns(), "payload table contains id + num_units"); + assertEquals(19999L, payload.getRowCount()); + AssertUtils.assertColumnsAreEqual(expectedIds, payload.getColumn(0), "id"); + } finally { + closeAll(filterCols); + closeAll(payloadCols); + } + } + } + // -------------------------------------------------------------------- // Tests: materializeAllColumns() // -------------------------------------------------------------------- @@ -825,6 +861,49 @@ void testMaterializePayloadColumnsChunkExactTotal(@TempDir Path tmp) throws IOEx } } + /** + * Verifies the chunked payload pipeline drains the same 19,999 rows when page pruning + * falls back to page-header row counts on a file with no page index. + */ + @Test + void testMaterializePayloadColumnsChunkPagePruningWithoutPageIndex(@TempDir Path tmp) + throws IOException { + try (OpenReader open = + OpenReader.multiPage(tmp).withFilter("zip_code", BinaryOperator.GREATER, 100000)) { + HybridScanReader reader = open.reader; + assertEquals(0L, reader.pageIndexByteRange().size(), + "Fixture must have no page index so the header-derived fallback is exercised"); + int[] survived = reader.filterRowGroupsWithStats(reader.allRowGroups()); + DeviceMemoryBuffer[] filterCols = copyRangesToDevice( + open.file, reader.filterColumnChunksByteRanges(survived)); + DeviceMemoryBuffer[] payloadCols = copyRangesToDevice( + open.file, reader.payloadColumnChunksByteRanges(survived)); + try { + reader.setupChunkingForFilterColumns(0L, 0L, survived, false, filterCols); + while (reader.hasNextTableChunk()) { + reader.materializeFilterColumnsChunk().close(); + } + try (ColumnVector rowMask = reader.takeFilterRowMask()) { + assertEquals(60000L, rowMask.getRowCount(), "Mask spans row group 1"); + assertEquals(19999L, countTrue(rowMask), "zip_code 100,001–119,999 survive"); + reader.setupChunkingForPayloadColumns(0L, 0L, survived, rowMask, true, payloadCols); + long total = 0; + while (reader.hasNextTableChunk()) { + try (Table chunk = reader.materializePayloadColumnsChunk(rowMask)) { + assertEquals(2, chunk.getNumberOfColumns()); + total += chunk.getRowCount(); + } + } + assertEquals(19999L, total, + "Header-derived page pruning must not drop or duplicate selected rows"); + } + } finally { + closeAll(filterCols); + closeAll(payloadCols); + } + } + } + // -------------------------------------------------------------------- // Tests: setupChunkingForAllColumns() / materializeAllColumnsChunk() // @@ -1212,9 +1291,19 @@ static OpenReader pageIndex(Path tmp) throws IOException { return openFromFile(pq, DEFAULT_COLS); } + /** A single 100-row group, small enough that every column chunk holds one data page. */ static OpenReader rowGroupStats(Path tmp) throws IOException { File pq = tmp.resolve("fixture.parquet").toFile(); - writeRowGroupStatsParquet(pq); + writeNoPageIndexParquet(pq, 100, 1, + ParquetWriterOptions.StatisticsFrequency.ROWGROUP); + return openFromFile(pq, DEFAULT_COLS); + } + + /** Two 60,000-row groups, so each column chunk spans 3 data pages (20,000 rows each). */ + static OpenReader multiPage(Path tmp) throws IOException { + File pq = tmp.resolve("fixture.parquet").toFile(); + writeNoPageIndexParquet(pq, 60_000, 2, + ParquetWriterOptions.StatisticsFrequency.PAGE); return openFromFile(pq, DEFAULT_COLS); } @@ -1323,28 +1412,34 @@ private static int writeFixtureParquet(File path) { } /** - * Writes a small Parquet file with {@code ROWGROUP}-level statistics: row-group min/max - * are recorded but no page index (no {@code ColumnIndex}/{@code OffsetIndex}) is emitted. - * Includes a low-cardinality {@code num_units} column ({1, 2, 3} cycle) so the writer's - * ADAPTIVE dictionary policy emits a dictionary; this lets tests exercise the - * "no page index, dict exists" path (see - * {@link #testSecondaryFiltersByteRangesEmptyForRowGroupStats}). + * Writes a Parquet file with no page index; only {@code COLUMN} statistics emit one. Both + * {@code id} and {@code zip_code} hold the globally sequential row index, and + * {@code num_units} cycles over {1, 2, 3}, low-cardinality enough that the ADAPTIVE + * dictionary policy emits a dictionary (see + * {@link #testSecondaryFiltersByteRangesEmptyForRowGroupStats}). The writer caps a page at + * 20,000 rows, so exceed that per group for chunks spanning several pages. */ - private static void writeRowGroupStatsParquet(File path) { - int rows = 100; + private static void writeNoPageIndexParquet(File path, int rowsPerGroup, int numGroups, + ParquetWriterOptions.StatisticsFrequency stats) { ParquetWriterOptions opts = ParquetWriterOptions.builder() .withNonNullableColumns("id", "zip_code", "num_units") - .withRowGroupSizeRows(rows) - .withStatisticsFrequency(ParquetWriterOptions.StatisticsFrequency.ROWGROUP) + .withRowGroupSizeRows(rowsPerGroup) + .withStatisticsFrequency(stats) .build(); - try (TableWriter writer = Table.writeParquetChunked(opts, path); - ColumnVector id = ColumnVector.fromInts(IntStream.range(0, rows).toArray()); - ColumnVector zipCode = ColumnVector.fromInts( - IntStream.range(0, rows).map(i -> 10000 + i).toArray()); - ColumnVector numUnits = ColumnVector.fromInts( - IntStream.range(0, rows).map(i -> 1 + (i % 3)).toArray()); - Table t = new Table(id, zipCode, numUnits)) { - writer.write(t); + try (TableWriter writer = Table.writeParquetChunked(opts, path)) { + for (int g = 0; g < numGroups; g++) { + int start = g * rowsPerGroup; + try (ColumnVector id = ColumnVector.fromInts( + IntStream.range(start, start + rowsPerGroup).toArray()); + ColumnVector zipCode = ColumnVector.fromInts( + IntStream.range(start, start + rowsPerGroup).toArray()); + ColumnVector numUnits = ColumnVector.fromInts( + IntStream.range(start, start + rowsPerGroup) + .map(i -> 1 + (i % 3)).toArray()); + Table t = new Table(id, zipCode, numUnits)) { + writer.write(t); + } + } } } diff --git a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py index 7bf3a19e1d13..aba3206012ff 100644 --- a/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py +++ b/python/pylibcudf/tests/io/test_experimental_hybrid_scan.py @@ -430,6 +430,73 @@ def test_hybrid_scan_materialize_columns( assert expected_arrow.equals(hybrid_arrow) +def test_hybrid_scan_payload_page_mask_without_page_index( + simple_parquet_bytes: bytes, + simple_hybrid_scan_reader: HybridScanReader, + simple_parquet_options: plc.io.parquet.ParquetReaderOptions, + simple_parquet_table: pa.Table, + num_rows: int, +) -> None: + """Test payload page pruning without a page index set up on the reader.""" + reader = simple_hybrid_scan_reader + row_groups = reader.all_row_groups(simple_parquet_options) + + # Keep the first half of the rows so the trailing data pages get pruned. + num_selected = num_rows // 2 + row_mask = plc.Column.from_arrow( + pa.array([i < num_selected for i in range(num_rows)], type=pa.bool_()) + ) + + payload_data = [ + plc.gpumemoryview( + rmm.DeviceBuffer.to_device( + simple_parquet_bytes[r.offset : r.offset + r.size], + plc.utils._get_stream(), + ) + ) + for r in reader.payload_column_chunks_byte_ranges( + row_groups, simple_parquet_options + ) + ] + synchronize_stream() + + # Chunks can disagree on field nullability, so compare row values only. + def to_rows(tbl: plc.Table) -> list: + return ( + tbl.to_arrow() + .rename_columns(simple_parquet_table.column_names) + .to_pylist() + ) + + expected_rows = simple_parquet_table.slice(0, num_selected).to_pylist() + + payload_result = reader.materialize_payload_columns( + row_groups, + payload_data, + row_mask, + UseDataPageMask.YES, + simple_parquet_options, + ) + synchronize_stream() + assert to_rows(payload_result.tbl) == expected_rows + + reader.setup_chunking_for_payload_columns( + 256, + 0, + row_groups, + row_mask, + UseDataPageMask.YES, + payload_data, + simple_parquet_options, + ) + chunked_rows = [] + while reader.has_next_table_chunk(): + chunk = reader.materialize_payload_columns_chunk(row_mask) + chunked_rows.extend(to_rows(chunk.tbl)) + synchronize_stream() + assert chunked_rows == expected_rows + + @pytest.mark.parametrize("stream", [None, Stream()]) def test_hybrid_scan_single_step_materialize( simple_parquet_bytes: bytes,