From 0a768580d1e9654b1128328d3fd51a17c15462df Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 01:21:56 +0000 Subject: [PATCH 01/23] Hybrid scan avoids null masks for REQUIRED unless page pruning --- .../experimental/hybrid_scan_helpers.cpp | 32 ------------ .../experimental/hybrid_scan_helpers.hpp | 31 +++++------ .../parquet/experimental/hybrid_scan_impl.cpp | 25 +++++++++ .../parquet/experimental/hybrid_scan_impl.hpp | 5 ++ .../hybrid_scan_multifile_test.cpp | 52 +++++++++++++++++++ 5 files changed, 95 insertions(+), 50 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 6e81d1bf3a93..fe7526d78ea9 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -143,38 +143,6 @@ aggregate_reader_metadata::aggregate_reader_metadata( initialize_internals(use_arrow_schema, has_cols_from_mismatched_srcs); } -void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, - bool has_cols_from_mismatched_srcs) -{ - keyval_maps = collect_keyval_metadata(); - schema_idx_maps = init_schema_idx_maps(has_cols_from_mismatched_srcs); - num_rows = calc_num_rows(); - num_row_groups = calc_num_row_groups(); - - // Force all non-nullable (REQUIRED) columns to be nullable without modifying REPEATED columns to - // preserve list structures - std::for_each(per_file_metadata.begin(), per_file_metadata.end(), [](auto& pfm) { - auto& schema = pfm.schema; - std::for_each(schema.begin() + 1, schema.end(), [](auto& col) { - // TODO: Store information of whichever column schema we modified here and restore it to - // `REQUIRED` if we end up not pruning any pages out of it - if (col.repetition_type == FieldRepetitionType::REQUIRED) { - col.repetition_type = FieldRepetitionType::OPTIONAL; - } - }); - }); - - // Collect and apply arrow:schema from Parquet's key value metadata section - if (use_arrow_schema) { - apply_arrow_schema(); - - // Erase ARROW_SCHEMA_KEY from the output pfm if exists - std::for_each(keyval_maps.begin(), keyval_maps.end(), [](auto& pfm) { - pfm.erase(cudf::io::parquet::detail::ARROW_SCHEMA_KEY); - }); - } -} - std::vector aggregate_reader_metadata::page_index_byte_ranges() const { std::vector page_index_byte_ranges; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index 59591d438913..334f47f6670e 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -78,19 +78,6 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { cuda::stream_ref stream) const; public: - /** - * @brief Check whether selected columns have column and offset indexes - * - * Schema indices are mapped to each source before locating the column chunks. - * - * @param row_group_indices Row group indices, one vector per source - * @param schema_indices Schema indices from the first source - * @return A pair indicating column-index and offset-index presence, respectively - */ - [[nodiscard]] std::pair page_index_presence( - std::span const> row_group_indices, - std::span schema_indices) const; - /** * @brief Constructor for aggregate_reader_metadata * @@ -118,11 +105,6 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { aggregate_reader_metadata(aggregate_reader_metadata&&) = default; aggregate_reader_metadata& operator=(aggregate_reader_metadata&&) = default; - /** - * @brief Initialize the internal variables - */ - void initialize_internals(bool use_arrow_schema, bool has_cols_from_mismatched_srcs); - /** * @brief Fetch the byte range of the page index in each Parquet file * @@ -137,6 +119,19 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { */ [[nodiscard]] std::vector parquet_metadatas() const; + /** + * @brief Check whether selected columns have column and offset indexes + * + * Schema indices are mapped to each source before locating the column chunks. + * + * @param row_group_indices Row group indices, one vector per source + * @param schema_indices Schema indices from the first source + * @return A pair indicating column-index and offset-index presence, respectively + */ + [[nodiscard]] std::pair page_index_presence( + std::span const> row_group_indices, + std::span schema_indices) const; + /** * @brief Setup and populate the page index structs in every source's `FileMetaData` * diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 82b15871408f..47ba2523e2d3 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -118,6 +118,25 @@ namespace { } // namespace +void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() +{ + if (std::all_of(_pass_page_mask.begin(), _pass_page_mask.end(), std::identity{})) { return; } + + auto const mark_buffers_nullable = [](auto const& self, std::span buffers) { + for (auto& buffer : buffers) { + // Page pruning synthesizes null rows at every nesting level except list elements. + if ((buffer.user_data & parquet::detail::PARQUET_COLUMN_BUFFER_FLAG_HAS_LIST_PARENT) == 0) { + buffer.is_nullable = true; + } + self(self, buffer.children); + } + }; + + // Mark both the output and template buffers nullable when page pruning synthesizes null rows + mark_buffers_nullable(mark_buffers_nullable, _output_buffers); + mark_buffers_nullable(mark_buffers_nullable, _output_buffers_template); +} + hybrid_scan_reader_impl::hybrid_scan_reader_impl( cudf::host_span const> footer_bytes, parquet_reader_options const& options) @@ -1439,6 +1458,9 @@ void hybrid_scan_reader_impl::set_pass_page_mask(std::span data_page // Make sure we inserted exactly the number of pages for this pass CUDF_EXPECTS(_pass_page_mask.size() == pass->pages.size(), "Encountered mismatch in number of pass pages and page mask size"); + + // Mark output buffers nullable when page pruning produces nulls + mark_buffers_nullable_for_pruned_pages(); } void hybrid_scan_reader_impl::set_sparse_pass_page_mask( @@ -1487,6 +1509,9 @@ void hybrid_scan_reader_impl::set_sparse_pass_page_mask( // Make sure we inserted exactly the number of pages for this pass. CUDF_EXPECTS(_pass_page_mask.size() == pass->pages.size(), "Encountered mismatch in number of pass pages and page mask size"); + + // Mark output buffers nullable when page pruning produces nulls + mark_buffers_nullable_for_pruned_pages(); } } // 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 b7c38ac71a69..6acb397c9a53 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -394,6 +394,11 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { */ void set_sparse_pass_page_mask(std::span const> page_data); + /** + * @brief Mark output buffers nullable when page pruning synthesizes null rows + */ + void mark_buffers_nullable_for_pruned_pages(); + /** * @brief Select the columns to be read based on the read mode * diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp index 47c8164468b6..974df68963f1 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp @@ -8,6 +8,7 @@ #include "tests/io/parquet_common.hpp" #include +#include #include #include @@ -472,3 +473,54 @@ TEST_F(HybridScanMultifileTest, SparsePayloadEmptyAndAllPrunedPageData) EXPECT_FALSE(reader->has_next_table_chunk()); } } +TEST_F(HybridScanMultifileTest, ChunkedAllColumnsPreservesRequiredNullability) +{ + // Use several pages so a nullable output would require a non-trivial validity allocation. + auto constexpr num_rows = 2 * page_size_for_ordered_tests; + auto values = cuda::counting_iterator{0}; + cudf::test::fixed_width_column_wrapper required_column(values, values + num_rows); + auto const input_table = cudf::table_view{{required_column}}; + + cudf::io::table_input_metadata metadata(input_table); + metadata.column_metadata[0].set_name("required"); + metadata.column_metadata[0].set_nullability(false); + + std::vector parquet_buffer; + auto const write_options = + cudf::io::parquet_writer_options::builder(cudf::io::sink_info{&parquet_buffer}, input_table) + .metadata(metadata) + .max_page_size_rows(page_size_for_ordered_tests) + .build(); + cudf::io::write_parquet(write_options); + + auto const source_info = build_source_info({parquet_buffer}); + auto const options = cudf::io::parquet_reader_options::builder().build(); + auto inputs = multifile_inputs(source_info); + auto reader = + cudf::io::parquet::experimental::hybrid_scan_multifile{inputs.footer_byte_spans, options}; + + ASSERT_EQ(reader.parquet_metadatas().front().schema[1].repetition_type, + cudf::io::parquet::FieldRepetitionType::REQUIRED); + + auto const row_groups = reader.all_row_groups(options); + auto const stream = cudf::get_default_stream(); + auto const mr = cudf::get_current_device_resource_ref(); + auto column_data = fetch_multisource_device_data( + inputs, reader.all_column_chunks_byte_ranges(row_groups, options), stream, mr); + reader.setup_chunking_for_all_columns( + 0, 0, row_groups, column_data.flat_spans, options, stream, mr); + + ASSERT_TRUE(reader.has_next_table_chunk()); + auto const result = reader.materialize_all_columns_chunk(); + EXPECT_FALSE(reader.has_next_table_chunk()); + + auto const result_column = result.tbl->view().column(0); + EXPECT_EQ(result_column.null_count(), 0); + EXPECT_FALSE(result_column.nullable()); + EXPECT_EQ(result_column.null_mask(), nullptr); + + auto const standard_result = cudf::io::read_parquet( + cudf::io::parquet_reader_options::builder(source_info).build(), stream, mr); + EXPECT_FALSE(standard_result.tbl->view().column(0).nullable()); + CUDF_TEST_EXPECT_TABLES_EQUAL(input_table, result.tbl->view()); +} From cd9563f3d27a22a6118deb451fbc77cc7bd3de8f Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 01:35:12 +0000 Subject: [PATCH 02/23] Only mark buffers with pruned pages --- .../parquet/experimental/hybrid_scan_impl.cpp | 36 +++++++++++++++---- 1 file changed, 30 insertions(+), 6 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 47ba2523e2d3..ac4eb6f77a76 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -27,6 +27,7 @@ #include #include #include +#include #include #include @@ -120,9 +121,24 @@ namespace { void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() { - if (std::all_of(_pass_page_mask.begin(), _pass_page_mask.end(), std::identity{})) { return; } - - auto const mark_buffers_nullable = [](auto const& self, std::span buffers) { + auto const& pass = *_pass_itm_data; + auto buffers_with_pruned_pages = std::vector(_output_buffers.size(), false); + + // Helper to get pruned page indices + auto pruned_page_indices = + std::views::iota(std::size_t{0}, _pass_page_mask.size()) | + std::views::filter([&](auto page_idx) { return not _pass_page_mask[page_idx]; }); + + // Mark buffers with pruned pages + std::ranges::for_each(pruned_page_indices, [&](auto page_idx) { + auto const& chunk = pass.chunks[pass.pages[page_idx].chunk_idx]; + auto const& input_column = _input_columns[chunk.src_col_index]; + buffers_with_pruned_pages[input_column.nesting.front()] = true; + }); + + // Helper to mark a buffer and its children nullable + auto const mark_buffers_nullable = [](auto const& self, + std::span buffers) -> void { for (auto& buffer : buffers) { // Page pruning synthesizes null rows at every nesting level except list elements. if ((buffer.user_data & parquet::detail::PARQUET_COLUMN_BUFFER_FLAG_HAS_LIST_PARENT) == 0) { @@ -132,9 +148,17 @@ void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() } }; - // Mark both the output and template buffers nullable when page pruning synthesizes null rows - mark_buffers_nullable(mark_buffers_nullable, _output_buffers); - mark_buffers_nullable(mark_buffers_nullable, _output_buffers_template); + // Mark buffers with pruned pages as nullable + auto buffers_with_pruned_page_indices = + std::views::iota(std::size_t{0}, _output_buffers.size()) | + std::views::filter([&](auto buffer_idx) { return buffers_with_pruned_pages[buffer_idx]; }); + std::ranges::for_each(buffers_with_pruned_page_indices, [&](auto buffer_idx) { + mark_buffers_nullable(mark_buffers_nullable, + std::span{&_output_buffers[buffer_idx], 1}); + mark_buffers_nullable( + mark_buffers_nullable, + std::span{&_output_buffers_template[buffer_idx], 1}); + }); } hybrid_scan_reader_impl::hybrid_scan_reader_impl( From 7167cc433aea0a0d3e5ba7c4e8d4746bee60824c Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 01:46:06 +0000 Subject: [PATCH 03/23] Simplify --- .../experimental/hybrid_scan_helpers.cpp | 40 +++++++++++-------- cpp/src/io/parquet/reader_impl_helpers.cpp | 2 + 2 files changed, 25 insertions(+), 17 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index fe7526d78ea9..86f612807753 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -36,6 +36,22 @@ using text::byte_range_info; namespace { +// Construct a vector of FileMetaData from the input footer bytes +[[nodiscard]] std::vector parquet_metadatas_from_footer_bytes( + cudf::host_span const> footer_bytes) +{ + std::vector parquet_metadatas; + parquet_metadatas.reserve(footer_bytes.size()); + std::transform(footer_bytes.begin(), + footer_bytes.end(), + std::back_inserter(parquet_metadatas), + [](auto const& footer_bytes) { + metadata parsed_metadata{footer_bytes}; + return FileMetaData{std::move(parsed_metadata)}; + }); + return parquet_metadatas; +} + // Construct a vector of all row group indices from the input vectors [[nodiscard]] auto all_row_group_indices( std::span const> row_group_indices) @@ -116,31 +132,21 @@ aggregate_reader_metadata::aggregate_reader_metadata( cudf::host_span const> footer_bytes, bool use_arrow_schema, bool has_cols_from_mismatched_srcs) - : aggregate_reader_metadata_base(host_span const>{}, false, false) + : aggregate_reader_metadata_base(parquet_metadatas_from_footer_bytes(footer_bytes), + use_arrow_schema, + has_cols_from_mismatched_srcs) { - CUDF_EXPECTS(not footer_bytes.empty(), "At least one source must be provided"); - per_file_metadata.reserve(footer_bytes.size()); - std::transform(footer_bytes.begin(), - footer_bytes.end(), - std::back_inserter(per_file_metadata), - [](auto const& fb) { return metadata{fb}; }); - initialize_internals(use_arrow_schema, has_cols_from_mismatched_srcs); } aggregate_reader_metadata::aggregate_reader_metadata( cudf::host_span parquet_metadatas, bool use_arrow_schema, bool has_cols_from_mismatched_srcs) - : aggregate_reader_metadata_base(host_span const>{}, false, false) + : aggregate_reader_metadata_base( + std::vector{parquet_metadatas.begin(), parquet_metadatas.end()}, + use_arrow_schema, + has_cols_from_mismatched_srcs) { - CUDF_EXPECTS(not parquet_metadatas.empty(), "At least one source must be provided"); - per_file_metadata.reserve(parquet_metadatas.size()); - // Just copy over the FileMetaData structs to the internal metadata structs - std::transform(parquet_metadatas.begin(), - parquet_metadatas.end(), - std::back_inserter(per_file_metadata), - [](auto const& parquet_metadata) { return metadata{parquet_metadata}; }); - initialize_internals(use_arrow_schema, has_cols_from_mismatched_srcs); } std::vector aggregate_reader_metadata::page_index_byte_ranges() const diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 5b51c2fd871e..6cf6609cc21b 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -956,6 +956,8 @@ aggregate_reader_metadata::aggregate_reader_metadata(std::vector&& bool use_arrow_schema, bool has_cols_from_mismatched_srcs) { + CUDF_EXPECTS(not parquet_metadatas.empty(), "At least one source must be provided"); + per_file_metadata.reserve(parquet_metadatas.size()); std::transform(std::make_move_iterator(parquet_metadatas.begin()), std::make_move_iterator(parquet_metadatas.end()), From 9ed5801c803989568c1efc0092278a1b159da03b Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 01:53:46 +0000 Subject: [PATCH 04/23] Minor fix --- cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp | 5 +++++ cpp/src/io/parquet/reader_impl_helpers.cpp | 2 -- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 86f612807753..6bf21950251e 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -40,6 +40,9 @@ namespace { [[nodiscard]] std::vector parquet_metadatas_from_footer_bytes( cudf::host_span const> footer_bytes) { + CUDF_EXPECTS( + not footer_bytes.empty(), "At least one source must be provided", std::invalid_argument); + std::vector parquet_metadatas; parquet_metadatas.reserve(footer_bytes.size()); std::transform(footer_bytes.begin(), @@ -147,6 +150,8 @@ aggregate_reader_metadata::aggregate_reader_metadata( use_arrow_schema, has_cols_from_mismatched_srcs) { + CUDF_EXPECTS( + not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); } std::vector aggregate_reader_metadata::page_index_byte_ranges() const diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 6cf6609cc21b..5b51c2fd871e 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -956,8 +956,6 @@ aggregate_reader_metadata::aggregate_reader_metadata(std::vector&& bool use_arrow_schema, bool has_cols_from_mismatched_srcs) { - CUDF_EXPECTS(not parquet_metadatas.empty(), "At least one source must be provided"); - per_file_metadata.reserve(parquet_metadatas.size()); std::transform(std::make_move_iterator(parquet_metadatas.begin()), std::make_move_iterator(parquet_metadatas.end()), From 1af5913433e36c9df0b4e3b7c0d16f058d0e5d85 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 02:06:33 +0000 Subject: [PATCH 05/23] Minor --- .../hybrid_scan_multifile_test.cpp | 51 +++++-------------- 1 file changed, 14 insertions(+), 37 deletions(-) diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp index 974df68963f1..b98ec3743b12 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp @@ -473,54 +473,31 @@ TEST_F(HybridScanMultifileTest, SparsePayloadEmptyAndAllPrunedPageData) EXPECT_FALSE(reader->has_next_table_chunk()); } } -TEST_F(HybridScanMultifileTest, ChunkedAllColumnsPreservesRequiredNullability) + +TEST_F(HybridScanMultifileTest, AllColumnsPreservesRequiredNullability) { - // Use several pages so a nullable output would require a non-trivial validity allocation. - auto constexpr num_rows = 2 * page_size_for_ordered_tests; - auto values = cuda::counting_iterator{0}; - cudf::test::fixed_width_column_wrapper required_column(values, values + num_rows); - auto const input_table = cudf::table_view{{required_column}}; - - cudf::io::table_input_metadata metadata(input_table); - metadata.column_metadata[0].set_name("required"); - metadata.column_metadata[0].set_nullability(false); - - std::vector parquet_buffer; - auto const write_options = - cudf::io::parquet_writer_options::builder(cudf::io::sink_info{&parquet_buffer}, input_table) - .metadata(metadata) - .max_page_size_rows(page_size_for_ordered_tests) - .build(); - cudf::io::write_parquet(write_options); - - auto const source_info = build_source_info({parquet_buffer}); - auto const options = cudf::io::parquet_reader_options::builder().build(); - auto inputs = multifile_inputs(source_info); + auto const [input_table, parquet_buffer] = create_parquet_with_stats(); + auto const parquet_buffers = std::vector>{parquet_buffer}; + auto const source_info = build_source_info(parquet_buffers); + auto const options = cudf::io::parquet_reader_options::builder().build(); + auto inputs = multifile_inputs(source_info); auto reader = cudf::io::parquet::experimental::hybrid_scan_multifile{inputs.footer_byte_spans, options}; + // Check footer metadata ASSERT_EQ(reader.parquet_metadatas().front().schema[1].repetition_type, cudf::io::parquet::FieldRepetitionType::REQUIRED); + // Materialize all columns auto const row_groups = reader.all_row_groups(options); auto const stream = cudf::get_default_stream(); auto const mr = cudf::get_current_device_resource_ref(); auto column_data = fetch_multisource_device_data( inputs, reader.all_column_chunks_byte_ranges(row_groups, options), stream, mr); - reader.setup_chunking_for_all_columns( - 0, 0, row_groups, column_data.flat_spans, options, stream, mr); - - ASSERT_TRUE(reader.has_next_table_chunk()); - auto const result = reader.materialize_all_columns_chunk(); - EXPECT_FALSE(reader.has_next_table_chunk()); - - auto const result_column = result.tbl->view().column(0); - EXPECT_EQ(result_column.null_count(), 0); - EXPECT_FALSE(result_column.nullable()); - EXPECT_EQ(result_column.null_mask(), nullptr); + auto const result = + reader.materialize_all_columns(row_groups, column_data.flat_spans, options, stream, mr); - auto const standard_result = cudf::io::read_parquet( - cudf::io::parquet_reader_options::builder(source_info).build(), stream, mr); - EXPECT_FALSE(standard_result.tbl->view().column(0).nullable()); - CUDF_TEST_EXPECT_TABLES_EQUAL(input_table, result.tbl->view()); + // Check results + EXPECT_FALSE(result.tbl->view().column(0).nullable()); + CUDF_TEST_EXPECT_TABLES_EQUAL(input_table->view(), result.tbl->view()); } From 0d06e6f03f83e54e05c237a05d7dec0987c71bd4 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 02:08:45 +0000 Subject: [PATCH 06/23] Minor --- cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 6bf21950251e..faba65f4f33a 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -40,9 +40,6 @@ namespace { [[nodiscard]] std::vector parquet_metadatas_from_footer_bytes( cudf::host_span const> footer_bytes) { - CUDF_EXPECTS( - not footer_bytes.empty(), "At least one source must be provided", std::invalid_argument); - std::vector parquet_metadatas; parquet_metadatas.reserve(footer_bytes.size()); std::transform(footer_bytes.begin(), @@ -139,6 +136,8 @@ aggregate_reader_metadata::aggregate_reader_metadata( use_arrow_schema, has_cols_from_mismatched_srcs) { + CUDF_EXPECTS( + not footer_bytes.empty(), "At least one source must be provided", std::invalid_argument); } aggregate_reader_metadata::aggregate_reader_metadata( From 3ce9db0bdcc1415de8d5af19401e45971dfdd111 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Tue, 25 Aug 2026 02:49:41 +0000 Subject: [PATCH 07/23] Hybrid scan constructor that move in footers --- .../cudf/io/experimental/hybrid_scan_multifile.hpp | 9 +++++++++ .../io/parquet/experimental/hybrid_scan_helpers.cpp | 13 +++++++++++++ .../io/parquet/experimental/hybrid_scan_helpers.hpp | 11 +++++++++++ .../io/parquet/experimental/hybrid_scan_impl.cpp | 10 ++++++++++ .../io/parquet/experimental/hybrid_scan_impl.hpp | 9 +++++++++ .../parquet/experimental/hybrid_scan_multifile.cpp | 6 ++++++ .../hybrid_scan_multifile_filters_test.cpp | 7 ++++--- 7 files changed, 62 insertions(+), 3 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index c75fa3d186d3..ac67b8c75cac 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -76,6 +76,15 @@ class hybrid_scan_multifile { explicit hybrid_scan_multifile(cudf::host_span parquet_metadata, parquet_reader_options const& options); + /** + * @brief Constructor that takes ownership of pre-populated Parquet file metadata + * + * @param parquet_metadata Pre-populated Parquet file metadata, one per source + * @param options Parquet reader options + */ + explicit hybrid_scan_multifile(std::vector&& parquet_metadata, + parquet_reader_options const& options); + /** * @brief Destructor for the multi-file experimental Parquet reader */ diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index faba65f4f33a..0ee5207d91ac 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -153,6 +153,19 @@ aggregate_reader_metadata::aggregate_reader_metadata( not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); } +aggregate_reader_metadata::aggregate_reader_metadata( + std::vector&& parquet_metadatas, + bool use_arrow_schema, + bool has_cols_from_mismatched_srcs) + : aggregate_reader_metadata_base( + std::move(parquet_metadatas), + use_arrow_schema, + has_cols_from_mismatched_srcs) +{ + CUDF_EXPECTS( + not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); +} + std::vector aggregate_reader_metadata::page_index_byte_ranges() const { std::vector page_index_byte_ranges; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index 334f47f6670e..7f798d6dbed4 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -100,6 +100,17 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { bool use_arrow_schema, bool has_cols_from_mismatched_srcs); + /** + * @brief Constructor that takes ownership of pre-populated Parquet file metadata + * + * @param parquet_metadatas Pre-populated Parquet file metadata, one per source + * @param use_arrow_schema Whether to use Arrow schema + * @param has_cols_from_mismatched_srcs Whether to have columns from mismatched sources + */ + aggregate_reader_metadata(std::vector&& parquet_metadatas, + bool use_arrow_schema, + bool has_cols_from_mismatched_srcs); + aggregate_reader_metadata(aggregate_reader_metadata const&) = delete; aggregate_reader_metadata& operator=(aggregate_reader_metadata const&) = delete; aggregate_reader_metadata(aggregate_reader_metadata&&) = default; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index ac4eb6f77a76..395601a0a061 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -181,6 +181,16 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl( _extended_metadata = static_cast(_metadata.get()); } +hybrid_scan_reader_impl::hybrid_scan_reader_impl(std::vector&& parquet_metadatas, + parquet_reader_options const& options) +{ + _metadata = + std::make_shared(std::move(parquet_metadatas), + options.is_enabled_use_arrow_schema(), + has_cols_from_mismatched_sources(options)); + _extended_metadata = static_cast(_metadata.get()); +} + hybrid_scan_reader_impl::hybrid_scan_reader_impl( std::shared_ptr metadata) { diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 6acb397c9a53..114de393edbb 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -58,6 +58,15 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { explicit hybrid_scan_reader_impl(cudf::host_span parquet_metadatas, parquet_reader_options const& options); + /** + * @brief Constructor that takes ownership of pre-populated Parquet file metadata + * + * @param parquet_metadatas Pre-populated Parquet file metadata, one per source + * @param options Parquet reader options + */ + explicit hybrid_scan_reader_impl(std::vector&& parquet_metadatas, + parquet_reader_options const& options); + /** * @brief Constructor that takes shared ownership of pre-parsed Parquet metadata * diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp index 259435401aab..493c86a8f3b2 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_multifile.cpp @@ -26,6 +26,12 @@ hybrid_scan_multifile::hybrid_scan_multifile(cudf::host_span { } +hybrid_scan_multifile::hybrid_scan_multifile(std::vector&& parquet_metadata, + parquet_reader_options const& options) + : _impl{std::make_unique(std::move(parquet_metadata), options)} +{ +} + hybrid_scan_multifile::~hybrid_scan_multifile() = default; std::vector hybrid_scan_multifile::parquet_metadatas() const 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 6e7cdad91bde..fbf28c85bb01 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_filters_test.cpp @@ -180,16 +180,17 @@ TEST_F(HybridScanMultifileFiltersTest, Metadata) EXPECT_EQ(reader->total_rows_in_row_groups(input_row_group_indices), 2 * rows_per_row_group * num_sources); - // Construct a new reader from a span of existing FileMetaData + // Move the existing FileMetaData into a new reader without copying it. + auto const num_row_groups = parquet_metadata.front().row_groups.size(); auto const reader_with_existing_metadata = std::make_unique( - cudf::host_span{parquet_metadata}, options); + std::move(parquet_metadata), options); // Check if the new metadata is the same as the existing one auto const new_metadata = reader_with_existing_metadata->parquet_metadatas(); ASSERT_EQ(new_metadata.size(), num_sources); EXPECT_TRUE(std::all_of(new_metadata.begin(), new_metadata.end(), [&](auto const& meta) { - return meta.row_groups.size() == parquet_metadata.front().row_groups.size(); + return meta.row_groups.size() == num_row_groups; })); } From 1e9df79d5bc3f0850440491904d5e272d2c1b56d Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 03:23:39 +0000 Subject: [PATCH 08/23] Improve --- .../parquet/experimental/hybrid_scan_impl.cpp | 41 +++++++++++-------- .../parquet/experimental/hybrid_scan_impl.hpp | 2 + 2 files changed, 25 insertions(+), 18 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index ac4eb6f77a76..4b554652b78c 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -121,20 +121,27 @@ namespace { void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() { - auto const& pass = *_pass_itm_data; - auto buffers_with_pruned_pages = std::vector(_output_buffers.size(), false); - - // Helper to get pruned page indices - auto pruned_page_indices = - std::views::iota(std::size_t{0}, _pass_page_mask.size()) | - std::views::filter([&](auto page_idx) { return not _pass_page_mask[page_idx]; }); - - // Mark buffers with pruned pages - std::ranges::for_each(pruned_page_indices, [&](auto page_idx) { - auto const& chunk = pass.chunks[pass.pages[page_idx].chunk_idx]; - auto const& input_column = _input_columns[chunk.src_col_index]; - buffers_with_pruned_pages[input_column.nesting.front()] = true; - }); + // Initialize flags for buffers with pruned pages + if (_buffers_with_pruned_pages.size() != _output_buffers.size()) { + _buffers_with_pruned_pages.resize(_output_buffers.size(), false); + } + + // Mark pruned buffers for this pass + if (_pass_itm_data) { + auto const& pass = *_pass_itm_data; + + // Helper to get pruned page indices + auto pruned_page_indices = + std::views::iota(std::size_t{0}, _pass_page_mask.size()) | + std::views::filter([&](auto page_idx) { return not _pass_page_mask[page_idx]; }); + + // Mark buffers with pruned pages + std::ranges::for_each(pruned_page_indices, [&](auto page_idx) { + auto const& chunk = pass.chunks[pass.pages[page_idx].chunk_idx]; + auto const& input_column = _input_columns[chunk.src_col_index]; + _buffers_with_pruned_pages[input_column.nesting.front()] = true; + }); + } // Helper to mark a buffer and its children nullable auto const mark_buffers_nullable = [](auto const& self, @@ -151,13 +158,10 @@ void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() // Mark buffers with pruned pages as nullable auto buffers_with_pruned_page_indices = std::views::iota(std::size_t{0}, _output_buffers.size()) | - std::views::filter([&](auto buffer_idx) { return buffers_with_pruned_pages[buffer_idx]; }); + std::views::filter([&](auto buffer_idx) { return _buffers_with_pruned_pages[buffer_idx]; }); std::ranges::for_each(buffers_with_pruned_page_indices, [&](auto buffer_idx) { mark_buffers_nullable(mark_buffers_nullable, std::span{&_output_buffers[buffer_idx], 1}); - mark_buffers_nullable( - mark_buffers_nullable, - std::span{&_output_buffers_template[buffer_idx], 1}); }); } @@ -1124,6 +1128,7 @@ void hybrid_scan_reader_impl::reset_internal_state() _pass_itm_data.reset(); _pass_page_mask.clear(); _subpass_page_mask.reset(); + _buffers_with_pruned_pages.clear(); _output_metadata.reset(); _sparse_page_io = false; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 6acb397c9a53..553c426ae08b 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -621,6 +621,8 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { std::optional> _filter_columns_names; + std::vector _buffers_with_pruned_pages; + cudf::size_type _row_mask_offset{0}; bool _output_chunk_produced{false}; From 75d714177681945f258ff8a17656dee3a4ea5e59 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb <14217455+mhaseeb123@users.noreply.github.com> Date: Tue, 25 Aug 2026 03:42:53 +0000 Subject: [PATCH 09/23] bug fixing --- .../parquet/experimental/hybrid_scan_impl.cpp | 72 +++++++++++-------- .../parquet/experimental/hybrid_scan_impl.hpp | 7 +- 2 files changed, 48 insertions(+), 31 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 4b554652b78c..9732c2dddf09 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -79,6 +79,24 @@ namespace { return output_dtypes; } +/** + * @brief Construct a vector of empty-like buffers from the input buffers + * + * @param buffers Input buffers + * @return Vector of empty-like buffers + */ +[[nodiscard]] std::vector make_empty_like_column_buffers( + std::span buffers) +{ + std::vector empty_buffers; + empty_buffers.reserve(buffers.size()); + std::transform( + buffers.begin(), buffers.end(), std::back_inserter(empty_buffers), [](auto const& buffer) { + return inline_column_buffer::empty_like(buffer); + }); + return empty_buffers; +} + /** * @brief Count the number of row groups in the input * @@ -121,27 +139,16 @@ namespace { void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() { - // Initialize flags for buffers with pruned pages - if (_buffers_with_pruned_pages.size() != _output_buffers.size()) { - _buffers_with_pruned_pages.resize(_output_buffers.size(), false); - } - - // Mark pruned buffers for this pass - if (_pass_itm_data) { - auto const& pass = *_pass_itm_data; - - // Helper to get pruned page indices - auto pruned_page_indices = - std::views::iota(std::size_t{0}, _pass_page_mask.size()) | - std::views::filter([&](auto page_idx) { return not _pass_page_mask[page_idx]; }); - - // Mark buffers with pruned pages - std::ranges::for_each(pruned_page_indices, [&](auto page_idx) { - auto const& chunk = pass.chunks[pass.pages[page_idx].chunk_idx]; - auto const& input_column = _input_columns[chunk.src_col_index]; - _buffers_with_pruned_pages[input_column.nesting.front()] = true; - }); - } + auto const& pass = *_pass_itm_data; + auto buffers_with_pruned_pages = std::vector(_output_buffers.size(), false); + auto pruned_page_indices = + std::views::iota(std::size_t{0}, _pass_page_mask.size()) | + std::views::filter([&](auto page_idx) { return not _pass_page_mask[page_idx]; }); + std::ranges::for_each(pruned_page_indices, [&](auto page_idx) { + auto const& chunk = pass.chunks[pass.pages[page_idx].chunk_idx]; + auto const& input_column = _input_columns[chunk.src_col_index]; + buffers_with_pruned_pages[input_column.nesting.front()] = true; + }); // Helper to mark a buffer and its children nullable auto const mark_buffers_nullable = [](auto const& self, @@ -158,10 +165,13 @@ void hybrid_scan_reader_impl::mark_buffers_nullable_for_pruned_pages() // Mark buffers with pruned pages as nullable auto buffers_with_pruned_page_indices = std::views::iota(std::size_t{0}, _output_buffers.size()) | - std::views::filter([&](auto buffer_idx) { return _buffers_with_pruned_pages[buffer_idx]; }); + std::views::filter([&](auto buffer_idx) { return buffers_with_pruned_pages[buffer_idx]; }); std::ranges::for_each(buffers_with_pruned_page_indices, [&](auto buffer_idx) { mark_buffers_nullable(mark_buffers_nullable, std::span{&_output_buffers[buffer_idx], 1}); + mark_buffers_nullable( + mark_buffers_nullable, + std::span{&_output_buffers_template[buffer_idx], 1}); }); } @@ -271,14 +281,16 @@ void hybrid_scan_reader_impl::select_columns(read_columns_mode read_columns_mode CUDF_EXPECTS(_input_columns.size() > 0 and _output_buffers.size() > 0, "No columns selected"); - // Clear the output buffers templates - _output_buffers_template.clear(); + // Save original output-buffer schema for reuse across materialization passes. + _original_output_buffers_template = make_empty_like_column_buffers(_output_buffers); + + // Initialize mutable output-buffer template for this materialization pass. + reset_output_buffers_template(); +} - // Save the states of the output buffers for reuse. - std::transform(_output_buffers.begin(), - _output_buffers.end(), - std::back_inserter(_output_buffers_template), - [](auto const& buff) { return inline_column_buffer::empty_like(buff); }); +void hybrid_scan_reader_impl::reset_output_buffers_template() +{ + _output_buffers_template = make_empty_like_column_buffers(_original_output_buffers_template); } std::vector> hybrid_scan_reader_impl::all_row_groups( @@ -323,6 +335,7 @@ void hybrid_scan_reader_impl::prepare_materialization(read_columns_mode read_col reset_internal_state(); initialize_options(options, num_sources, stream, mr); select_columns(read_columns_mode, options); + reset_output_buffers_template(); } std::vector> @@ -1128,7 +1141,6 @@ void hybrid_scan_reader_impl::reset_internal_state() _pass_itm_data.reset(); _pass_page_mask.clear(); _subpass_page_mask.reset(); - _buffers_with_pruned_pages.clear(); _output_metadata.reset(); _sparse_page_io = false; diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 553c426ae08b..26d2ba82f3e5 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -399,6 +399,11 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { */ void mark_buffers_nullable_for_pruned_pages(); + /** + * @brief Initialize the mutable output-buffer template for this materialization + */ + void reset_output_buffers_template(); + /** * @brief Select the columns to be read based on the read mode * @@ -621,7 +626,7 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { std::optional> _filter_columns_names; - std::vector _buffers_with_pruned_pages; + std::vector _original_output_buffers_template; cudf::size_type _row_mask_offset{0}; bool _output_chunk_produced{false}; From 8da4ba9279db4b10e82355aa44fdc25d598c0eb3 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Tue, 25 Aug 2026 04:00:09 +0000 Subject: [PATCH 10/23] style fix --- .../io/parquet/experimental/hybrid_scan_helpers.cpp | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index 0ee5207d91ac..adb38b61c120 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -153,14 +153,11 @@ aggregate_reader_metadata::aggregate_reader_metadata( not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); } -aggregate_reader_metadata::aggregate_reader_metadata( - std::vector&& parquet_metadatas, - bool use_arrow_schema, - bool has_cols_from_mismatched_srcs) +aggregate_reader_metadata::aggregate_reader_metadata(std::vector&& parquet_metadatas, + bool use_arrow_schema, + bool has_cols_from_mismatched_srcs) : aggregate_reader_metadata_base( - std::move(parquet_metadatas), - use_arrow_schema, - has_cols_from_mismatched_srcs) + std::move(parquet_metadatas), use_arrow_schema, has_cols_from_mismatched_srcs) { CUDF_EXPECTS( not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); From 3594c05577be0a089159ba11d8fdfcb99255ad06 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Tue, 25 Aug 2026 21:00:58 +0000 Subject: [PATCH 11/23] Avoid multiple page index setups --- cpp/src/io/parquet/reader_impl_helpers.cpp | 16 +++++++++++++++- cpp/src/io/parquet/reader_impl_helpers.hpp | 3 +++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 5b51c2fd871e..2c2dbbeff506 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -509,7 +509,17 @@ void metadata::sanitize_schema() process(0); } -metadata::metadata(FileMetaData&& other) : FileMetaData(std::move(other)) {} +metadata::metadata(FileMetaData&& other) : FileMetaData(std::move(other)) +{ + // Since page index is set up for all or no row groups, just check if any column chunk has it set. + // Update this check if this behavior changes in the future. + is_page_index_setup = + std::any_of(row_groups.cbegin(), row_groups.cend(), [](auto const& row_group) { + return std::any_of(row_group.columns.cbegin(), row_group.columns.cend(), [](auto const& col) { + return col.column_index.has_value() or col.offset_index.has_value(); + }); + }); +} metadata::metadata(datasource* source, bool read_page_indexes) { @@ -549,6 +559,8 @@ metadata::metadata(datasource* source, bool read_page_indexes) void metadata::setup_page_index(cudf::host_span page_index_bytes, int64_t min_offset) { + if (is_page_index_setup) { return; } + CUDF_FUNC_RANGE(); // Flatten all columns into a single vector for easier task distribution @@ -624,6 +636,8 @@ void metadata::setup_page_index(cudf::host_span page_index_bytes, read_column_indexes(cp, col_ref.get()); } } + + is_page_index_setup = true; } metadata::~metadata() diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index 6e9eb5e18266..52744ee36120 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -171,6 +171,9 @@ struct metadata : public FileMetaData { protected: void sanitize_schema(); + + private: + bool is_page_index_setup = false; }; /** From dc9b533baeb867ad3df352e5d7ded65c374b3a34 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 00:06:44 +0000 Subject: [PATCH 12/23] Address comments from @vuule --- .../io/experimental/hybrid_scan_multifile.hpp | 6 + .../io/parquet/experimental/hybrid_scan.cpp | 4 +- .../experimental/hybrid_scan_helpers.cpp | 32 +---- .../experimental/hybrid_scan_helpers.hpp | 6 + .../parquet/experimental/hybrid_scan_impl.cpp | 17 +-- .../parquet/experimental/hybrid_scan_impl.hpp | 4 +- cpp/src/io/parquet/reader_impl.hpp | 27 ++-- cpp/src/io/parquet/reader_impl_helpers.cpp | 119 ++++++++++-------- cpp/src/io/parquet/reader_impl_helpers.hpp | 39 ++++++ .../hybrid_scan_multifile_test.cpp | 97 ++++++++++++++ 10 files changed, 251 insertions(+), 100 deletions(-) diff --git a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp index ac67b8c75cac..23e4b5b44129 100644 --- a/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp +++ b/cpp/include/cudf/io/experimental/hybrid_scan_multifile.hpp @@ -61,6 +61,8 @@ class hybrid_scan_multifile { /** * @brief Constructor for the multi-file experimental Parquet reader * + * @throws std::invalid_argument if no sources are provided + * * @param footer_bytes Host span of Parquet file footer byte spans, one per source * @param options Parquet reader options */ @@ -70,6 +72,8 @@ class hybrid_scan_multifile { /** * @brief Constructor for the multi-file experimental Parquet reader * + * @throws std::invalid_argument if no sources are provided + * * @param parquet_metadata Host span of pre-populated Parquet file metadata, one per source * @param options Parquet reader options */ @@ -79,6 +83,8 @@ class hybrid_scan_multifile { /** * @brief Constructor that takes ownership of pre-populated Parquet file metadata * + * @throws std::invalid_argument if no sources are provided + * * @param parquet_metadata Pre-populated Parquet file metadata, one per source * @param options Parquet reader options */ diff --git a/cpp/src/io/parquet/experimental/hybrid_scan.cpp b/cpp/src/io/parquet/experimental/hybrid_scan.cpp index 09faac8c261a..879eec607960 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan.cpp @@ -18,7 +18,7 @@ hybrid_scan_metadata::hybrid_scan_metadata(cudf::host_span footer : _metadata{std::make_shared( std::vector>{footer_bytes}, options.is_enabled_use_arrow_schema(), - options.get_column_names().has_value() and options.is_enabled_allow_mismatched_pq_schemas())} + options.is_enabled_allow_mismatched_pq_schemas())} { } @@ -27,7 +27,7 @@ hybrid_scan_metadata::hybrid_scan_metadata(FileMetaData const& parquet_metadata, : _metadata{std::make_shared( std::vector{parquet_metadata}, options.is_enabled_use_arrow_schema(), - options.get_column_names().has_value() and options.is_enabled_allow_mismatched_pq_schemas())} + options.is_enabled_allow_mismatched_pq_schemas())} { } diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp index adb38b61c120..d0523b55fe71 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.cpp @@ -36,22 +36,6 @@ using text::byte_range_info; namespace { -// Construct a vector of FileMetaData from the input footer bytes -[[nodiscard]] std::vector parquet_metadatas_from_footer_bytes( - cudf::host_span const> footer_bytes) -{ - std::vector parquet_metadatas; - parquet_metadatas.reserve(footer_bytes.size()); - std::transform(footer_bytes.begin(), - footer_bytes.end(), - std::back_inserter(parquet_metadatas), - [](auto const& footer_bytes) { - metadata parsed_metadata{footer_bytes}; - return FileMetaData{std::move(parsed_metadata)}; - }); - return parquet_metadatas; -} - // Construct a vector of all row group indices from the input vectors [[nodiscard]] auto all_row_group_indices( std::span const> row_group_indices) @@ -132,25 +116,23 @@ aggregate_reader_metadata::aggregate_reader_metadata( cudf::host_span const> footer_bytes, bool use_arrow_schema, bool has_cols_from_mismatched_srcs) - : aggregate_reader_metadata_base(parquet_metadatas_from_footer_bytes(footer_bytes), - use_arrow_schema, - has_cols_from_mismatched_srcs) + : aggregate_reader_metadata( + parquet::detail::parallel_construct_metadatas( + footer_bytes, [](auto const& bytes) { return FileMetaData{metadata{bytes}}; }), + use_arrow_schema, + has_cols_from_mismatched_srcs) { - CUDF_EXPECTS( - not footer_bytes.empty(), "At least one source must be provided", std::invalid_argument); } aggregate_reader_metadata::aggregate_reader_metadata( cudf::host_span parquet_metadatas, bool use_arrow_schema, bool has_cols_from_mismatched_srcs) - : aggregate_reader_metadata_base( + : aggregate_reader_metadata( std::vector{parquet_metadatas.begin(), parquet_metadatas.end()}, use_arrow_schema, has_cols_from_mismatched_srcs) { - CUDF_EXPECTS( - not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); } aggregate_reader_metadata::aggregate_reader_metadata(std::vector&& parquet_metadatas, @@ -159,8 +141,6 @@ aggregate_reader_metadata::aggregate_reader_metadata(std::vector&& : aggregate_reader_metadata_base( std::move(parquet_metadatas), use_arrow_schema, has_cols_from_mismatched_srcs) { - CUDF_EXPECTS( - not parquet_metadatas.empty(), "At least one source must be provided", std::invalid_argument); } std::vector aggregate_reader_metadata::page_index_byte_ranges() const diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp index 7f798d6dbed4..72118214ade5 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_helpers.hpp @@ -81,6 +81,8 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { /** * @brief Constructor for aggregate_reader_metadata * + * @throws std::invalid_argument if no sources are provided + * * @param footer_bytes Host span of Parquet file footer buffer bytes, one per source * @param use_arrow_schema Whether to use Arrow schema * @param has_cols_from_mismatched_srcs Whether to have columns from mismatched sources @@ -92,6 +94,8 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { /** * @brief Constructor for aggregate_reader_metadata * + * @throws std::invalid_argument if no sources are provided + * * @param parquet_metadatas Host span of pre-populated Parquet file metadata, one per source * @param use_arrow_schema Whether to use Arrow schema * @param has_cols_from_mismatched_srcs Whether to have columns from mismatched sources @@ -103,6 +107,8 @@ class aggregate_reader_metadata : public aggregate_reader_metadata_base { /** * @brief Constructor that takes ownership of pre-populated Parquet file metadata * + * @throws std::invalid_argument if no sources are provided + * * @param parquet_metadatas Pre-populated Parquet file metadata, one per source * @param use_arrow_schema Whether to use Arrow schema * @param has_cols_from_mismatched_srcs Whether to have columns from mismatched sources diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index 1577c600f1b3..ba377c38cb75 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -180,7 +180,9 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl( parquet_reader_options const& options) { _metadata = std::make_shared( - footer_bytes, options.is_enabled_use_arrow_schema(), has_cols_from_mismatched_sources(options)); + footer_bytes, + options.is_enabled_use_arrow_schema(), + options.is_enabled_allow_mismatched_pq_schemas()); _extended_metadata = static_cast(_metadata.get()); } @@ -191,7 +193,7 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl( _metadata = std::make_shared(parquet_metadatas, options.is_enabled_use_arrow_schema(), - has_cols_from_mismatched_sources(options)); + options.is_enabled_allow_mismatched_pq_schemas()); _extended_metadata = static_cast(_metadata.get()); } @@ -201,7 +203,7 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl(std::vector&& par _metadata = std::make_shared(std::move(parquet_metadatas), options.is_enabled_use_arrow_schema(), - has_cols_from_mismatched_sources(options)); + options.is_enabled_allow_mismatched_pq_schemas()); _extended_metadata = static_cast(_metadata.get()); } @@ -294,12 +296,13 @@ void hybrid_scan_reader_impl::select_columns(read_columns_mode read_columns_mode // Save original output-buffer schema for reuse across materialization passes. _original_output_buffers_template = make_empty_like_column_buffers(_output_buffers); - // Initialize mutable output-buffer template for this materialization pass. - reset_output_buffers_template(); + // Initialize mutable output buffers for this materialization pass. + reset_output_buffers(); } -void hybrid_scan_reader_impl::reset_output_buffers_template() +void hybrid_scan_reader_impl::reset_output_buffers() { + _output_buffers = make_empty_like_column_buffers(_original_output_buffers_template); _output_buffers_template = make_empty_like_column_buffers(_original_output_buffers_template); } @@ -345,7 +348,7 @@ void hybrid_scan_reader_impl::prepare_materialization(read_columns_mode read_col reset_internal_state(); initialize_options(options, num_sources, stream, mr); select_columns(read_columns_mode, options); - reset_output_buffers_template(); + reset_output_buffers(); } std::vector> diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp index 2194bcb69a21..a2e2bda3c5f7 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.hpp @@ -409,9 +409,9 @@ class hybrid_scan_reader_impl : public parquet::detail::reader_impl { void mark_buffers_nullable_for_pruned_pages(); /** - * @brief Initialize the mutable output-buffer template for this materialization + * @brief Reset the output buffers and their template from the original selected-columns schema */ - void reset_output_buffers_template(); + void reset_output_buffers(); /** * @brief Select the columns to be read based on the read mode diff --git a/cpp/src/io/parquet/reader_impl.hpp b/cpp/src/io/parquet/reader_impl.hpp index e811971d0583..30d7550264f0 100644 --- a/cpp/src/io/parquet/reader_impl.hpp +++ b/cpp/src/io/parquet/reader_impl.hpp @@ -413,19 +413,6 @@ class reader_impl { _file_itm_data._current_input_pass < _file_itm_data.num_passes(); } - /** - * @brief Check if the user has specified columns from mismatched sources - * - * @param options Reader options - * @return True if the user has specified columns from mismatched sources - */ - [[nodiscard]] bool has_cols_from_mismatched_sources(parquet_reader_options const& options) const - { - return (options.get_column_names().has_value() or - options.get_column_field_ids().has_value()) and - options.is_enabled_allow_mismatched_pq_schemas(); - } - /** * @brief Effective `ignore_missing_columns` policy for column selection * @@ -440,6 +427,20 @@ class reader_impl { not(has_cols_from_mismatched_sources(options) and _metadata->get_num_sources() > 1); } + private: + /** + * @brief Check if the user has specified columns from mismatched sources + * + * @param options Reader options + * @return True if the user has specified columns from mismatched sources + */ + [[nodiscard]] bool has_cols_from_mismatched_sources(parquet_reader_options const& options) const + { + return (options.get_column_names().has_value() or + options.get_column_field_ids().has_value()) and + options.is_enabled_allow_mismatched_pq_schemas(); + } + protected: /** * @brief Check if the user has specified custom row bounds diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 2c2dbbeff506..ad9fe8dcee67 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -289,6 +289,39 @@ struct schema_child_lookup { std::unordered_map schema_idx_caches; }; +/** + * @brief Mark a column and its children in `schema` optional if it is optional in `other` + * + * Columns are matched by name as mismatched schemas may order their columns differently. + Columns missing from `other` are skipped as they are validated at column selection stage. + * + * @param schema Schema of the first source, updated in place + * @param other Schema of another source + * @param schema_idx Index of the current element in `schema` + * @param other_idx Index of the matching element in `other` + */ +void propagate_optional_fields(std::vector& schema, + std::vector const& other, + size_type schema_idx, + size_type other_idx) +{ + if (schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED and + other[other_idx].repetition_type != FieldRepetitionType::REQUIRED) { + schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; + } + + auto const& other_children = other[other_idx].children_idx; + for (auto const child_idx : schema[schema_idx].children_idx) { + auto const other_child_iter = + std::find_if(other_children.begin(), other_children.end(), [&](auto const idx) { + return other[idx].name == schema[child_idx].name; + }); + if (other_child_iter != other_children.end()) { + propagate_optional_fields(schema, other, child_idx, *other_child_iter); + } + } +} + } // namespace type_id to_type_id(SchemaElement const& schema, @@ -658,26 +691,9 @@ metadata::~metadata() std::vector aggregate_reader_metadata::metadatas_from_sources( host_span const> sources, bool read_page_indexes) { - // Avoid using the thread pool for a single source - if (sources.size() == 1) { - std::vector result; - result.emplace_back(sources[0].get(), read_page_indexes); - return result; - } - - std::vector> metadata_ctor_tasks; - metadata_ctor_tasks.reserve(sources.size()); - for (auto const& source : sources) { - metadata_ctor_tasks.emplace_back(cudf::detail::host_worker_pool().submit_task( - [source = source.get(), read_page_indexes] { return metadata{source, read_page_indexes}; })); - } - std::vector metadatas; - metadatas.reserve(sources.size()); - std::transform(metadata_ctor_tasks.begin(), - metadata_ctor_tasks.end(), - std::back_inserter(metadatas), - [](std::future& task) { return std::move(task).get(); }); - return metadatas; + return parallel_construct_metadatas(sources, [read_page_indexes](auto const& source) { + return metadata{source.get(), read_page_indexes}; + }); } std::vector> @@ -944,18 +960,9 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, // Mark the column schema in the first (default) source as nullable if it is nullable in any of // the input sources. This avoids recomputing this within build_column() and // populate_metadata(). - std::for_each( - cuda::counting_iterator{static_cast(1)}, - cuda::counting_iterator{schema.size()}, - [&](auto const schema_idx) { - if (schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED and - std::any_of( - per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { - return pfm.schema[schema_idx].repetition_type != FieldRepetitionType::REQUIRED; - })) { - schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; - } - }); + std::for_each(per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { + propagate_optional_fields(schema, pfm.schema, 0, 0); + }); } // Collect and apply arrow:schema from Parquet's key value metadata section @@ -970,6 +977,10 @@ aggregate_reader_metadata::aggregate_reader_metadata(std::vector&& bool use_arrow_schema, bool has_cols_from_mismatched_srcs) { + CUDF_EXPECTS(not parquet_metadatas.empty(), + "Encountered an empty vector of parquet metadatas (sources)", + std::invalid_argument); + per_file_metadata.reserve(parquet_metadatas.size()); std::transform(std::make_move_iterator(parquet_metadatas.begin()), std::make_move_iterator(parquet_metadatas.end()), @@ -2204,6 +2215,30 @@ aggregate_reader_metadata::select_columns( } }; + // Maps a top-level column's schema_idx across the rest of the data sources if we are reading from + // mismatched Parquet sources. `col_name_info` is null when all of the column's children are + // selected. + auto map_column_across_sources = [&](column_name_info const* col_name_info, + std::string const& col_name, + int const src_schema_idx) { + if (per_file_metadata.size() == 1 or schema_idx_maps.empty()) { return; } + + auto constexpr root_idx = 0; + std::for_each( + cuda::counting_iterator{static_cast(1)}, + cuda::counting_iterator{per_file_metadata.size()}, + [&](auto const src_idx) { + // Ensure that each top level column exists in the destination schema tree. + auto const dst_schema_idx = + schema_lookup.find_target_schema_child(root_idx, root_idx, col_name, src_idx); + CUDF_EXPECTS( + dst_schema_idx != -1, + std::format("Encountered missing top-level column '{}' across Parquet sources", col_name), + std::invalid_argument); + map_column(col_name_info, src_schema_idx, dst_schema_idx, src_idx); + }); + }; + std::vector output_column_schemas; // @@ -2230,6 +2265,7 @@ aggregate_reader_metadata::select_columns( for (auto const& schema_idx : root.children_idx) { build_column(nullptr, schema_idx, output_columns, false); output_column_schemas.push_back(schema_idx); + map_column_across_sources(nullptr, get_schema(schema_idx).name, schema_idx); } } else { struct path_info { @@ -2350,24 +2386,7 @@ aggregate_reader_metadata::select_columns( bool const valid_column = build_column(&col, top_level_col_schema_idx, output_columns, false); if (valid_column) { output_column_schemas.push_back(top_level_col_schema_idx); - - // Map the column's schema_idx across the rest of the data sources if required. - if (per_file_metadata.size() > 1 and not schema_idx_maps.empty()) { - std::for_each( - cuda::counting_iterator{static_cast(1)}, - cuda::counting_iterator{per_file_metadata.size()}, - [&](auto const src_idx) { - // Ensure that each top level column exists in the destination schema tree. - auto const dst_col_schema_idx = - schema_lookup.find_target_schema_child(root_idx, root_idx, col.name, src_idx); - CUDF_EXPECTS( - dst_col_schema_idx != -1, - std::format("Encountered missing top-level column '{}' across Parquet sources", - col.name), - std::invalid_argument); - map_column(&col, top_level_col_schema_idx, dst_col_schema_idx, src_idx); - }); - } + map_column_across_sources(&col, col.name, top_level_col_schema_idx); } } } diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index 52744ee36120..8d499c1fed55 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -7,26 +7,65 @@ #include "parquet_gpu.hpp" +#include #include #include #include #include #include +#include +#include #include #include +#include #include #include #include #include #include #include +#include #include #include #include namespace cudf::io::parquet::detail { +/** + * @brief Construct metadatas from inputs using the host worker pool for multiple inputs + * + * @param inputs Metadata construction inputs, one per source + * @param op Operation constructing a metadata object from one input + * @return Constructed metadata objects, in input order + */ +template +[[nodiscard]] auto parallel_construct_metadatas(cudf::host_span inputs, UnaryOp op) +{ + using result_type = std::invoke_result_t; + + std::vector results; + results.reserve(inputs.size()); + + // Avoid using the thread pool for a single input + if (inputs.size() == 1) { + results.emplace_back(op(inputs.front())); + return results; + } + + std::vector> tasks; + tasks.reserve(inputs.size()); + std::transform(inputs.begin(), inputs.end(), std::back_inserter(tasks), [&op](T const& input) { + return cudf::detail::host_worker_pool().submit_task( + [&op, input_ptr = &input] { return op(*input_ptr); }); + }); + + std::transform( + tasks.begin(), tasks.end(), std::back_inserter(results), [](auto& task) { return task.get(); }); + + return results; +} + /** * @brief page location and size info */ diff --git a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp index b98ec3743b12..28f5264250b0 100644 --- a/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_multifile_test.cpp @@ -501,3 +501,100 @@ TEST_F(HybridScanMultifileTest, AllColumnsPreservesRequiredNullability) EXPECT_FALSE(result.tbl->view().column(0).nullable()); CUDF_TEST_EXPECT_TABLES_EQUAL(input_table->view(), result.tbl->view()); } + +TEST_F(HybridScanMultifileTest, ReadColumnsFromMismatchedSchemas) +{ + // Create two sources with mismatched schemas + auto const buffer_a = std::get<1>(create_parquet_with_stats()); + auto const buffer_b = std::get<1>(create_parquet_with_stats( + 100, cudf::io::compression_type::AUTO, {"col2", "col0", "col1"}, {2, 0, 1})); + + auto const parquet_buffers = std::vector>{buffer_a, buffer_b}; + auto const source_info = build_source_info(parquet_buffers); + auto inputs = multifile_inputs(source_info); + auto const stream = cudf::get_default_stream(); + auto const mr = cudf::get_current_device_resource_ref(); + + // Reading mismatched schemas must be opted into, even without a column projection + EXPECT_THROW(cudf::io::parquet::experimental::hybrid_scan_multifile( + inputs.footer_byte_spans, cudf::io::parquet_reader_options::builder().build()), + cudf::logic_error); + + // Expected table from the regular reader + auto const expected = + cudf::io::read_parquet(cudf::io::parquet_reader_options::builder(source_info) + .allow_mismatched_pq_schemas(true) + .column_names({"col0", "col1", "col2"}) + .build(), + stream, + mr); + + auto options = + cudf::io::parquet_reader_options::builder().allow_mismatched_pq_schemas(true).build(); + auto const reader = + cudf::io::parquet::experimental::hybrid_scan_multifile{inputs.footer_byte_spans, options}; + auto const row_groups = reader.all_row_groups(options); + + // Single step materialize with hybrid scan + { + auto column_data = fetch_multisource_device_data( + inputs, reader.all_column_chunks_byte_ranges(row_groups, options), stream, mr); + auto const result = + reader.materialize_all_columns(row_groups, column_data.flat_spans, options, stream, mr); + CUDF_TEST_EXPECT_TABLES_EQUAL(expected.tbl->view(), result.tbl->view()); + } + + // Two step materialize with hybrid scan + { + auto literal_value = cudf::numeric_scalar(std::numeric_limits::min()); + auto literal = cudf::ast::literal(literal_value); + auto col_ref = cudf::ast::column_name_reference("col0"); + auto filter = cudf::ast::operation(cudf::ast::ast_operator::GREATER_EQUAL, col_ref, literal); + + options.set_filter(filter); + reader.reset_column_selection(); + + auto row_mask = reader.build_all_true_row_mask(row_groups, stream, mr); + auto row_mask_view = row_mask->mutable_view(); + + auto filter_column_chunks = fetch_multisource_device_data( + inputs, reader.filter_column_chunks_byte_ranges(row_groups, options), stream, mr); + auto const filter_result = reader.materialize_filter_columns(row_groups, + filter_column_chunks.flat_spans, + row_mask_view, + use_data_page_mask::NO, + options, + stream, + mr); + + auto payload_column_chunks = fetch_multisource_device_data( + inputs, reader.payload_column_chunks_byte_ranges(row_groups, options), stream, mr); + auto const payload_result = reader.materialize_payload_columns(row_groups, + payload_column_chunks.flat_spans, + row_mask_view, + use_data_page_mask::NO, + options, + stream, + mr); + + CUDF_TEST_EXPECT_TABLES_EQUAL(expected.tbl->select({0}), filter_result.tbl->view()); + CUDF_TEST_EXPECT_TABLES_EQUAL(expected.tbl->select({1, 2}), payload_result.tbl->view()); + } +} + +TEST_F(HybridScanMultifileTest, EmptySources) +{ + // Arrow schema is applied during metadata construction, so make sure empty inputs are rejected + // before any metadata is touched + auto const options = cudf::io::parquet_reader_options::builder().use_arrow_schema(true).build(); + + EXPECT_THROW(cudf::io::parquet::experimental::hybrid_scan_multifile( + cudf::host_span const>{}, options), + std::invalid_argument); + EXPECT_THROW(cudf::io::parquet::experimental::hybrid_scan_multifile( + cudf::host_span{}, options), + std::invalid_argument); + EXPECT_THROW(cudf::io::parquet::experimental::hybrid_scan_multifile( + std::vector{}, options), + std::invalid_argument); +} From 8dbdfb32e4096c17d42788c7fe7b6a2f3d0b095d Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 00:27:52 +0000 Subject: [PATCH 13/23] Simplify --- cpp/src/io/parquet/reader_impl_chunking.cu | 54 +++++++++++---------- cpp/src/io/parquet/reader_impl_helpers.cpp | 55 ++++++++++------------ 2 files changed, 56 insertions(+), 53 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_chunking.cu b/cpp/src/io/parquet/reader_impl_chunking.cu index 1cb95bd35a97..f728ebfc00cc 100644 --- a/cpp/src/io/parquet/reader_impl_chunking.cu +++ b/cpp/src/io/parquet/reader_impl_chunking.cu @@ -19,6 +19,7 @@ #include #include +#include namespace cudf::io::parquet::detail { @@ -419,29 +420,32 @@ void reader_impl::create_global_chunk_info() auto const num_chunks = row_groups_info.size() * num_input_columns; // Mapping of input column to page index column - std::vector column_mapping; - - if (_has_offset_index and not row_groups_info.empty()) { - // use first row group to define mappings (assumes same schema for each file) - auto const& rg = row_groups_info[0]; - auto const& columns = _metadata->get_row_group(rg.index, rg.source_index).columns; - column_mapping.resize(num_input_columns); - std::transform( - _input_columns.begin(), _input_columns.end(), column_mapping.begin(), [&](auto const& col) { - // translate schema_idx into something we can use for the page indexes - if (auto it = std::find_if(columns.begin(), - columns.end(), - [&](auto const& col_chunk) { - return col_chunk.schema_idx == - _metadata->map_schema_index(col.schema_idx, - rg.source_index); - }); - it != columns.end()) { - return std::distance(columns.begin(), it); - } - CUDF_FAIL("cannot find column mapping"); - }); - } + auto column_mappings = std::unordered_map>{}; + + auto const column_mapping_for_source = [&](auto const& rg) -> std::vector const& { + auto const [iter, inserted] = column_mappings.try_emplace(rg.source_index); + if (inserted) { + auto const& columns = _metadata->get_row_group(rg.index, rg.source_index).columns; + auto& mapping = iter->second; + mapping.resize(num_input_columns); + std::transform( + _input_columns.begin(), _input_columns.end(), mapping.begin(), [&](auto const& col) { + // translate schema_idx into something we can use for the page indexes + if (auto it = std::find_if(columns.begin(), + columns.end(), + [&](auto const& col_chunk) { + return col_chunk.schema_idx == + _metadata->map_schema_index(col.schema_idx, + rg.source_index); + }); + it != columns.end()) { + return static_cast(std::distance(columns.begin(), it)); + } + CUDF_FAIL("cannot find column mapping"); + }); + } + return iter->second; + }; // Initialize column chunk information auto remaining_rows = num_rows; @@ -454,6 +458,8 @@ void reader_impl::create_global_chunk_info() auto row_group_rows = std::min(remaining_rows + adjusted_row_group_rows, row_group.num_rows); + auto const* const column_mapping = _has_offset_index ? &column_mapping_for_source(rg) : nullptr; + // generate ColumnChunkDesc objects for everything to be decoded (all input columns) for (size_t i = 0; i < num_input_columns; ++i) { auto col = _input_columns[i]; @@ -479,7 +485,7 @@ void reader_impl::create_global_chunk_info() // grab the column_chunk_info for each chunk (if it exists) column_chunk_info const* const chunk_info = - _has_offset_index ? &rg.column_chunks.value()[column_mapping[i]] : nullptr; + _has_offset_index ? &rg.column_chunks.value()[(*column_mapping)[i]] : nullptr; chunks.emplace_back(col_meta.total_compressed_size, nullptr, diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index ad9fe8dcee67..e7ba480c5d27 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -290,36 +290,20 @@ struct schema_child_lookup { }; /** - * @brief Mark a column and its children in `schema` optional if it is optional in `other` + * @brief Maps the dot-separated path of each schema element to its schema index * - * Columns are matched by name as mismatched schemas may order their columns differently. - Columns missing from `other` are skipped as they are validated at column selection stage. - * - * @param schema Schema of the first source, updated in place - * @param other Schema of another source - * @param schema_idx Index of the current element in `schema` - * @param other_idx Index of the matching element in `other` + * @param schema Schema of a source + * @return Map of schema element paths to their schema indices */ -void propagate_optional_fields(std::vector& schema, - std::vector const& other, - size_type schema_idx, - size_type other_idx) +[[nodiscard]] std::unordered_map map_schema_paths_to_indices( + std::vector const& schema) { - if (schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED and - other[other_idx].repetition_type != FieldRepetitionType::REQUIRED) { - schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; - } - - auto const& other_children = other[other_idx].children_idx; - for (auto const child_idx : schema[schema_idx].children_idx) { - auto const other_child_iter = - std::find_if(other_children.begin(), other_children.end(), [&](auto const idx) { - return other[idx].name == schema[child_idx].name; - }); - if (other_child_iter != other_children.end()) { - propagate_optional_fields(schema, other, child_idx, *other_child_iter); - } + auto schema_indices = std::unordered_map{}; + schema_indices.reserve(schema.size()); + for (auto schema_idx = size_type{1}; std::cmp_less(schema_idx, schema.size()); ++schema_idx) { + schema_indices.emplace(column_path_from_index(schema, schema_idx), schema_idx); } + return schema_indices; } } // namespace @@ -958,10 +942,23 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, } // Mark the column schema in the first (default) source as nullable if it is nullable in any of - // the input sources. This avoids recomputing this within build_column() and - // populate_metadata(). + // the input sources. Fields are matched by their path as mismatched sources may order their + // columns differently. Fields missing from the first source skipped and instead validated at + // column selection stage. + auto const schema_indices = map_schema_paths_to_indices(schema); std::for_each(per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { - propagate_optional_fields(schema, pfm.schema, 0, 0); + std::for_each( + cuda::counting_iterator{1}, + cuda::counting_iterator{static_cast(pfm.schema.size())}, + [&](auto const other_idx) { + if (pfm.schema[other_idx].repetition_type == FieldRepetitionType::REQUIRED) { return; } + auto const schema_idx_iter = + schema_indices.find(column_path_from_index(pfm.schema, other_idx)); + if (schema_idx_iter != schema_indices.end() and + schema[schema_idx_iter->second].repetition_type == FieldRepetitionType::REQUIRED) { + schema[schema_idx_iter->second].repetition_type = FieldRepetitionType::OPTIONAL; + } + }); }); } From 80dde76f56bb4ef2ae05df94396e1758427e4545 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 00:40:30 +0000 Subject: [PATCH 14/23] Simplify --- cpp/src/io/parquet/reader_impl_helpers.cpp | 43 +++++++++++----------- 1 file changed, 22 insertions(+), 21 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index e7ba480c5d27..ad65344bee54 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -300,9 +300,11 @@ struct schema_child_lookup { { auto schema_indices = std::unordered_map{}; schema_indices.reserve(schema.size()); - for (auto schema_idx = size_type{1}; std::cmp_less(schema_idx, schema.size()); ++schema_idx) { - schema_indices.emplace(column_path_from_index(schema, schema_idx), schema_idx); - } + std::for_each(cuda::counting_iterator{1}, + cuda::counting_iterator{static_cast(schema.size())}, + [&](auto const schema_idx) { + schema_indices.emplace(column_path_from_index(schema, schema_idx), schema_idx); + }); return schema_indices; } @@ -941,25 +943,24 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, } } - // Mark the column schema in the first (default) source as nullable if it is nullable in any of - // the input sources. Fields are matched by their path as mismatched sources may order their - // columns differently. Fields missing from the first source skipped and instead validated at - // column selection stage. + // Mark a field in the first source schema as nullable if it is nullable in any other source auto const schema_indices = map_schema_paths_to_indices(schema); - std::for_each(per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { - std::for_each( - cuda::counting_iterator{1}, - cuda::counting_iterator{static_cast(pfm.schema.size())}, - [&](auto const other_idx) { - if (pfm.schema[other_idx].repetition_type == FieldRepetitionType::REQUIRED) { return; } - auto const schema_idx_iter = - schema_indices.find(column_path_from_index(pfm.schema, other_idx)); - if (schema_idx_iter != schema_indices.end() and - schema[schema_idx_iter->second].repetition_type == FieldRepetitionType::REQUIRED) { - schema[schema_idx_iter->second].repetition_type = FieldRepetitionType::OPTIONAL; - } - }); - }); + for (auto const& pfm : per_file_metadata) { + std::for_each(cuda::counting_iterator{1}, + cuda::counting_iterator{static_cast(pfm.schema.size())}, + [&](auto const schema_idx) { + if (pfm.schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED) { + return; + } + auto const iter = + schema_indices.find(column_path_from_index(pfm.schema, schema_idx)); + if (iter == schema_indices.end()) { return; } + if (auto& repetition_type = schema[iter->second].repetition_type; + repetition_type == FieldRepetitionType::REQUIRED) { + repetition_type = FieldRepetitionType::OPTIONAL; + } + }); + } } // Collect and apply arrow:schema from Parquet's key value metadata section From 5347f61868d48b66b3ec0f9f4bccade732653300 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 00:49:12 +0000 Subject: [PATCH 15/23] Apply suggestion from @coderabbitai --- cpp/src/io/parquet/reader_impl_helpers.hpp | 26 +++++++++++++++++----- 1 file changed, 20 insertions(+), 6 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.hpp b/cpp/src/io/parquet/reader_impl_helpers.hpp index 8d499c1fed55..f5dafae6f5ab 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.hpp +++ b/cpp/src/io/parquet/reader_impl_helpers.hpp @@ -17,6 +17,7 @@ #include #include +#include #include #include #include @@ -55,13 +56,26 @@ template std::vector> tasks; tasks.reserve(inputs.size()); - std::transform(inputs.begin(), inputs.end(), std::back_inserter(tasks), [&op](T const& input) { - return cudf::detail::host_worker_pool().submit_task( - [&op, input_ptr = &input] { return op(*input_ptr); }); - }); - std::transform( - tasks.begin(), tasks.end(), std::back_inserter(results), [](auto& task) { return task.get(); }); + auto pending_exception = std::exception_ptr{}; + try { + std::transform(inputs.begin(), inputs.end(), std::back_inserter(tasks), [&op](T const& input) { + return cudf::detail::host_worker_pool().submit_task( + [&op, input_ptr = &input] { return op(*input_ptr); }); + }); + } catch (...) { + pending_exception = std::current_exception(); + } + + for (auto& task : tasks) { + try { + results.emplace_back(task.get()); + } catch (...) { + if (not pending_exception) { pending_exception = std::current_exception(); } + } + } + + if (pending_exception) { std::rethrow_exception(pending_exception); } return results; } From 64a41377034a5a85ad462cf09f9bf35f7294d03b Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 01:09:55 +0000 Subject: [PATCH 16/23] Minor --- cpp/src/io/parquet/reader_impl_helpers.cpp | 5 ++++- cpp/tests/io/parquet_reader_test.cpp | 21 +++++++++++++++++++++ 2 files changed, 25 insertions(+), 1 deletion(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index ad65344bee54..90f8e3c73298 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -303,7 +303,10 @@ struct schema_child_lookup { std::for_each(cuda::counting_iterator{1}, cuda::counting_iterator{static_cast(schema.size())}, [&](auto const schema_idx) { - schema_indices.emplace(column_path_from_index(schema, schema_idx), schema_idx); + auto const path = column_path_from_index(schema, schema_idx); + auto const [_, inserted] = schema_indices.emplace(path, schema_idx); + CUDF_EXPECTS( + inserted, "Ambiguous parquet schema path: " + path, std::invalid_argument); }); return schema_indices; } diff --git a/cpp/tests/io/parquet_reader_test.cpp b/cpp/tests/io/parquet_reader_test.cpp index 6e55ccbd5ca0..c9fe769bae12 100644 --- a/cpp/tests/io/parquet_reader_test.cpp +++ b/cpp/tests/io/parquet_reader_test.cpp @@ -6291,3 +6291,24 @@ TEST_F(ParquetReaderTest, NestedMismatchedSchemaColumnValidation) EXPECT_THROW(cudf::io::read_parquet(opts), std::invalid_argument); } } + +TEST_F(ParquetReaderTest, DuplicateDottedSchemaPaths) +{ + auto child = cudf::test::fixed_width_column_wrapper{1, 2, 3}; + auto nested = cudf::test::structs_column_wrapper{{child}}; + auto scalar = cudf::test::fixed_width_column_wrapper{4, 5, 6}; + auto const table = cudf::table_view{{nested, scalar}}; + + auto const path = temp_env->get_temp_filepath("DuplicateDottedSchemaPaths.parquet"); + auto metadata = cudf::io::table_input_metadata(table); + metadata.column_metadata[0].set_name("a"); + metadata.column_metadata[0].child(0).set_name("b"); + metadata.column_metadata[1].set_name("a.b"); + cudf::io::write_parquet( + cudf::io::parquet_writer_options::builder(cudf::io::sink_info{path}, table) + .metadata(std::move(metadata))); + + auto const options = + cudf::io::parquet_reader_options::builder(cudf::io::source_info{{path, path}}).build(); + EXPECT_THROW(cudf::io::read_parquet(options), std::invalid_argument); +} \ No newline at end of file From 46588d37371ca9b9fb6775129690694def793c21 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 01:41:07 +0000 Subject: [PATCH 17/23] Fix --- cpp/src/io/parquet/reader_impl_helpers.cpp | 102 +++++++++++---------- cpp/tests/io/parquet_reader_test.cpp | 45 ++++++++- 2 files changed, 96 insertions(+), 51 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 90f8e3c73298..0d58f536457f 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -289,28 +289,6 @@ struct schema_child_lookup { std::unordered_map schema_idx_caches; }; -/** - * @brief Maps the dot-separated path of each schema element to its schema index - * - * @param schema Schema of a source - * @return Map of schema element paths to their schema indices - */ -[[nodiscard]] std::unordered_map map_schema_paths_to_indices( - std::vector const& schema) -{ - auto schema_indices = std::unordered_map{}; - schema_indices.reserve(schema.size()); - std::for_each(cuda::counting_iterator{1}, - cuda::counting_iterator{static_cast(schema.size())}, - [&](auto const schema_idx) { - auto const path = column_path_from_index(schema, schema_idx); - auto const [_, inserted] = schema_indices.emplace(path, schema_idx); - CUDF_EXPECTS( - inserted, "Ambiguous parquet schema path: " + path, std::invalid_argument); - }); - return schema_indices; -} - } // namespace type_id to_type_id(SchemaElement const& schema, @@ -944,25 +922,21 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, } CUDF_EXPECTS(schema == pfm.schema, "All sources must have the same schema"); } - } - // Mark a field in the first source schema as nullable if it is nullable in any other source - auto const schema_indices = map_schema_paths_to_indices(schema); - for (auto const& pfm : per_file_metadata) { - std::for_each(cuda::counting_iterator{1}, - cuda::counting_iterator{static_cast(pfm.schema.size())}, - [&](auto const schema_idx) { - if (pfm.schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED) { - return; - } - auto const iter = - schema_indices.find(column_path_from_index(pfm.schema, schema_idx)); - if (iter == schema_indices.end()) { return; } - if (auto& repetition_type = schema[iter->second].repetition_type; - repetition_type == FieldRepetitionType::REQUIRED) { - repetition_type = FieldRepetitionType::OPTIONAL; - } - }); + // Mark a field in the first source's schema as nullable if it is nullable in any other + // source. + std::for_each( + cuda::counting_iterator{static_cast(1)}, + cuda::counting_iterator{schema.size()}, + [&](auto const schema_idx) { + if (schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED and + std::any_of( + per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { + return pfm.schema[schema_idx].repetition_type != FieldRepetitionType::REQUIRED; + })) { + schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; + } + }); } } @@ -1007,6 +981,10 @@ aggregate_reader_metadata::aggregate_reader_metadata( num_rows(calc_num_rows()), num_row_groups(calc_num_row_groups()) { + CUDF_EXPECTS(not per_file_metadata.empty(), + "Encountered an empty vector of parquet sources", + std::invalid_argument); + initialize_internals(use_arrow_schema, has_cols_from_mismatched_srcs); } @@ -2240,6 +2218,25 @@ aggregate_reader_metadata::select_columns( }); }; + // Marks a mapped field in the first source's schema as nullable if it is nullable in any other + // source. Must run after all columns have been mapped across sources and before build_column() + auto propagate_mapped_optional_fields = [&]() { + if (schema_idx_maps.empty()) { return; } + auto& schema = per_file_metadata.front().schema; + std::for_each( + cuda::counting_iterator{static_cast(1)}, + cuda::counting_iterator{per_file_metadata.size()}, + [&](auto const src_idx) { + auto const& other_schema = per_file_metadata[src_idx].schema; + for (auto const& [schema_idx, other_idx] : schema_idx_maps[src_idx - 1]) { + if (other_schema[other_idx].repetition_type != FieldRepetitionType::REQUIRED and + schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED) { + schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; + } + } + }); + }; + std::vector output_column_schemas; // @@ -2263,10 +2260,14 @@ aggregate_reader_metadata::select_columns( // auto const& root = get_schema(0); if (not use_names.has_value()) { + for (auto const& schema_idx : root.children_idx) { + map_column_across_sources(nullptr, get_schema(schema_idx).name, schema_idx); + } + propagate_mapped_optional_fields(); + for (auto const& schema_idx : root.children_idx) { build_column(nullptr, schema_idx, output_columns, false); output_column_schemas.push_back(schema_idx); - map_column_across_sources(nullptr, get_schema(schema_idx).name, schema_idx); } } else { struct path_info { @@ -2380,16 +2381,25 @@ aggregate_reader_metadata::select_columns( } } } + + // Map the column's schema_idx across the rest of the data sources and propagate nullability. + auto constexpr root_idx = 0; for (auto& col : selected_columns) { - auto constexpr root_idx = 0; - auto const& top_level_col_schema_idx = + auto const top_level_col_schema_idx = schema_lookup.find_schema_child_by_name(root_idx, col.name); - bool const valid_column = build_column(&col, top_level_col_schema_idx, output_columns, false); - if (valid_column) { - output_column_schemas.push_back(top_level_col_schema_idx); + if (top_level_col_schema_idx != -1) { map_column_across_sources(&col, col.name, top_level_col_schema_idx); } } + propagate_mapped_optional_fields(); + + for (auto& col : selected_columns) { + auto const top_level_col_schema_idx = + schema_lookup.find_schema_child_by_name(root_idx, col.name); + if (build_column(&col, top_level_col_schema_idx, output_columns, false)) { + output_column_schemas.push_back(top_level_col_schema_idx); + } + } } return std::make_tuple( diff --git a/cpp/tests/io/parquet_reader_test.cpp b/cpp/tests/io/parquet_reader_test.cpp index c9fe769bae12..21a966f144b5 100644 --- a/cpp/tests/io/parquet_reader_test.cpp +++ b/cpp/tests/io/parquet_reader_test.cpp @@ -688,23 +688,24 @@ TEST_F(ParquetReaderTest, SelectMismatchedStructChildByFieldId) .metadata(std::move(metadata_a)); cudf::io::write_parquet(write_args_a); - auto y_b = cudf::test::fixed_width_column_wrapper{40, 50}; + auto y_b = cudf::test::fixed_width_column_wrapper{{40, 50}, {true, false}}; auto x_b = cudf::test::fixed_width_column_wrapper{4, 5}; auto struct_b = cudf::test::structs_column_wrapper{{y_b, x_b}, {true, true}}.release(); cudf::table_view const table_b{{*struct_b}}; auto path_b = temp_env->get_temp_filepath("SelectNestedFieldIdChildOrderB.parquet"); cudf::io::table_input_metadata metadata_b(table_b); - metadata_b.column_metadata[0].set_name("record").set_parquet_field_id(1); - metadata_b.column_metadata[0].child(0).set_name("y").set_parquet_field_id(3); - metadata_b.column_metadata[0].child(1).set_name("x").set_parquet_field_id(2); + metadata_b.column_metadata[0].set_name("renamed_record").set_parquet_field_id(1); + metadata_b.column_metadata[0].child(0).set_name("renamed_y").set_parquet_field_id(3); + metadata_b.column_metadata[0].child(1).set_name("renamed_x").set_parquet_field_id(2); auto write_args_b = cudf::io::parquet_writer_options::builder(cudf::io::sink_info{path_b}, table_b) .metadata(std::move(metadata_b)); cudf::io::write_parquet(write_args_b); auto expected_x = cudf::test::fixed_width_column_wrapper{1, 2, 3, 4, 5}; - auto expected_y = cudf::test::fixed_width_column_wrapper{10, 20, 30, 40, 50}; + auto expected_y = cudf::test::fixed_width_column_wrapper{ + {10, 20, 30, 40, 50}, {true, true, true, true, false}}; auto expected_struct = cudf::test::structs_column_wrapper{{expected_x, expected_y}, {true, true, true, true, true}} .release(); @@ -5411,6 +5412,17 @@ TEST_F(ParquetReaderTest, LateBindSourceInfo) CUDF_TEST_EXPECT_TABLES_EQUAL(result.tbl->view(), expected->view()); } +TEST_F(ParquetReaderTest, EmptySourcesWithArrowSchema) +{ + auto sources = std::vector>{}; + auto file_metadatas = std::vector{}; + auto const options = cudf::io::parquet_reader_options::builder(cudf::io::source_info{}) + .use_arrow_schema(true) + .build(); + EXPECT_THROW(cudf::io::read_parquet(std::move(sources), std::move(file_metadatas), options), + std::invalid_argument); +} + TEST_F(ParquetReaderTest, InvalidFooterMagic) { auto const expected = create_random_fixed_table(4, 4, false); @@ -6311,4 +6323,27 @@ TEST_F(ParquetReaderTest, DuplicateDottedSchemaPaths) auto const options = cudf::io::parquet_reader_options::builder(cudf::io::source_info{{path, path}}).build(); EXPECT_THROW(cudf::io::read_parquet(options), std::invalid_argument); +} + +TEST_F(ParquetReaderTest, CaseInsensitiveMismatchedSchemasPropagateNullability) +{ + auto const required = cudf::test::fixed_width_column_wrapper{1, 2, 3}; + auto const optional = cudf::test::fixed_width_column_wrapper{{4, 5}, {true, false}}; + auto const required_path = + write_parquet_temp_file(cudf::table_view{{required}}, "CaseRequired.parquet", {"column"}); + auto const optional_path = + write_parquet_temp_file(cudf::table_view{{optional}}, "CaseOptional.parquet", {"COLUMN"}); + + auto const options = + cudf::io::parquet_reader_options::builder(cudf::io::source_info{{required_path, optional_path}}) + .allow_mismatched_pq_schemas(true) + .case_sensitive_names(false) + .column_names({"column"}) + .build(); + + // A non-nullable column in the first source but nullable in another must be read as nullable. + auto const expected = cudf::test::fixed_width_column_wrapper{ + {1, 2, 3, 4, 5}, {true, true, true, true, false}}; + auto const result = cudf::io::read_parquet(options); + CUDF_TEST_EXPECT_COLUMNS_EQUAL(result.tbl->view().column(0), expected); } \ No newline at end of file From 5cc15d6bf352f6d4167ff4134facfc24aaba3f66 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 02:04:22 +0000 Subject: [PATCH 18/23] Simplify --- cpp/src/io/parquet/reader_impl_helpers.cpp | 59 ++++++++-------------- 1 file changed, 22 insertions(+), 37 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 0d58f536457f..5db43b61eba8 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -289,6 +289,20 @@ struct schema_child_lookup { std::unordered_map schema_idx_caches; }; +/** + * @brief Marks a field in the destination schema as optional if it is optional in the source schema + * + * @param dst Destination schema element + * @param src Source schema element + */ +void propagate_optional_field(SchemaElement& dst, SchemaElement const& src) +{ + if (dst.repetition_type == FieldRepetitionType::REQUIRED and + src.repetition_type != FieldRepetitionType::REQUIRED) { + dst.repetition_type = FieldRepetitionType::OPTIONAL; + } +} + } // namespace type_id to_type_id(SchemaElement const& schema, @@ -924,18 +938,15 @@ void aggregate_reader_metadata::initialize_internals(bool use_arrow_schema, } // Mark a field in the first source's schema as nullable if it is nullable in any other - // source. + // source std::for_each( cuda::counting_iterator{static_cast(1)}, cuda::counting_iterator{schema.size()}, [&](auto const schema_idx) { - if (schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED and - std::any_of( - per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { - return pfm.schema[schema_idx].repetition_type != FieldRepetitionType::REQUIRED; - })) { - schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; - } + std::for_each( + per_file_metadata.begin() + 1, per_file_metadata.end(), [&](auto const& pfm) { + propagate_optional_field(schema[schema_idx], pfm.schema[schema_idx]); + }); }); } } @@ -2130,6 +2141,9 @@ aggregate_reader_metadata::select_columns( // Map the schema index from 0th tree (src) to the one in the current (dst) tree. schema_idx_map[src_schema_idx] = dst_schema_idx; + // Mark the field nullable in the first schema if it is nullable in the current tree. + propagate_optional_field(per_file_metadata.front().schema[src_schema_idx], dst_schema_elem); + // If src_schema_elem is a stub, it does not exist in the column_name_info and column_buffer // hierarchy. So continue on with mapping. if (src_schema_elem.is_stub()) { @@ -2218,25 +2232,6 @@ aggregate_reader_metadata::select_columns( }); }; - // Marks a mapped field in the first source's schema as nullable if it is nullable in any other - // source. Must run after all columns have been mapped across sources and before build_column() - auto propagate_mapped_optional_fields = [&]() { - if (schema_idx_maps.empty()) { return; } - auto& schema = per_file_metadata.front().schema; - std::for_each( - cuda::counting_iterator{static_cast(1)}, - cuda::counting_iterator{per_file_metadata.size()}, - [&](auto const src_idx) { - auto const& other_schema = per_file_metadata[src_idx].schema; - for (auto const& [schema_idx, other_idx] : schema_idx_maps[src_idx - 1]) { - if (other_schema[other_idx].repetition_type != FieldRepetitionType::REQUIRED and - schema[schema_idx].repetition_type == FieldRepetitionType::REQUIRED) { - schema[schema_idx].repetition_type = FieldRepetitionType::OPTIONAL; - } - } - }); - }; - std::vector output_column_schemas; // @@ -2262,10 +2257,6 @@ aggregate_reader_metadata::select_columns( if (not use_names.has_value()) { for (auto const& schema_idx : root.children_idx) { map_column_across_sources(nullptr, get_schema(schema_idx).name, schema_idx); - } - propagate_mapped_optional_fields(); - - for (auto const& schema_idx : root.children_idx) { build_column(nullptr, schema_idx, output_columns, false); output_column_schemas.push_back(schema_idx); } @@ -2390,12 +2381,6 @@ aggregate_reader_metadata::select_columns( if (top_level_col_schema_idx != -1) { map_column_across_sources(&col, col.name, top_level_col_schema_idx); } - } - propagate_mapped_optional_fields(); - - for (auto& col : selected_columns) { - auto const top_level_col_schema_idx = - schema_lookup.find_schema_child_by_name(root_idx, col.name); if (build_column(&col, top_level_col_schema_idx, output_columns, false)) { output_column_schemas.push_back(top_level_col_schema_idx); } From afb797d32452a7f2c4faa3c44e64ddbf3d332889 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 02:12:44 +0000 Subject: [PATCH 19/23] Minor simplfication --- cpp/src/io/parquet/reader_impl_helpers.cpp | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 5db43b61eba8..1199ec5ab5a3 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -2378,9 +2378,7 @@ aggregate_reader_metadata::select_columns( for (auto& col : selected_columns) { auto const top_level_col_schema_idx = schema_lookup.find_schema_child_by_name(root_idx, col.name); - if (top_level_col_schema_idx != -1) { - map_column_across_sources(&col, col.name, top_level_col_schema_idx); - } + map_column_across_sources(&col, col.name, top_level_col_schema_idx); if (build_column(&col, top_level_col_schema_idx, output_columns, false)) { output_column_schemas.push_back(top_level_col_schema_idx); } From 7a78a0cff1a1a29af782372f867ec0605cc51597 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 02:16:13 +0000 Subject: [PATCH 20/23] Minor --- cpp/src/io/parquet/reader_impl_helpers.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index 1199ec5ab5a3..a302476f6ed9 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -1981,6 +1981,8 @@ aggregate_reader_metadata::select_columns( auto const case_sensitive_names = selection_options.case_sensitive_names; auto const selection_mode = selection_options.selection_mode; + auto constexpr root_idx = 0; + // Setup schema lookup helper auto schema_lookup = schema_child_lookup{[&](int const schema_idx, int const src_idx) -> SchemaElement const& { @@ -2216,7 +2218,6 @@ aggregate_reader_metadata::select_columns( int const src_schema_idx) { if (per_file_metadata.size() == 1 or schema_idx_maps.empty()) { return; } - auto constexpr root_idx = 0; std::for_each( cuda::counting_iterator{static_cast(1)}, cuda::counting_iterator{per_file_metadata.size()}, @@ -2374,7 +2375,6 @@ aggregate_reader_metadata::select_columns( } // Map the column's schema_idx across the rest of the data sources and propagate nullability. - auto constexpr root_idx = 0; for (auto& col : selected_columns) { auto const top_level_col_schema_idx = schema_lookup.find_schema_child_by_name(root_idx, col.name); From afa5b5749c834fd07359899f4437d99ec284d7d5 Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 02:23:51 +0000 Subject: [PATCH 21/23] Remove stupid tests --- cpp/tests/io/parquet_reader_test.cpp | 21 --------------------- 1 file changed, 21 deletions(-) diff --git a/cpp/tests/io/parquet_reader_test.cpp b/cpp/tests/io/parquet_reader_test.cpp index 21a966f144b5..6f86cd4b5e69 100644 --- a/cpp/tests/io/parquet_reader_test.cpp +++ b/cpp/tests/io/parquet_reader_test.cpp @@ -6304,27 +6304,6 @@ TEST_F(ParquetReaderTest, NestedMismatchedSchemaColumnValidation) } } -TEST_F(ParquetReaderTest, DuplicateDottedSchemaPaths) -{ - auto child = cudf::test::fixed_width_column_wrapper{1, 2, 3}; - auto nested = cudf::test::structs_column_wrapper{{child}}; - auto scalar = cudf::test::fixed_width_column_wrapper{4, 5, 6}; - auto const table = cudf::table_view{{nested, scalar}}; - - auto const path = temp_env->get_temp_filepath("DuplicateDottedSchemaPaths.parquet"); - auto metadata = cudf::io::table_input_metadata(table); - metadata.column_metadata[0].set_name("a"); - metadata.column_metadata[0].child(0).set_name("b"); - metadata.column_metadata[1].set_name("a.b"); - cudf::io::write_parquet( - cudf::io::parquet_writer_options::builder(cudf::io::sink_info{path}, table) - .metadata(std::move(metadata))); - - auto const options = - cudf::io::parquet_reader_options::builder(cudf::io::source_info{{path, path}}).build(); - EXPECT_THROW(cudf::io::read_parquet(options), std::invalid_argument); -} - TEST_F(ParquetReaderTest, CaseInsensitiveMismatchedSchemasPropagateNullability) { auto const required = cudf::test::fixed_width_column_wrapper{1, 2, 3}; From e8b384fbb27aa2f52af2e5060e554bf5b3b70a9c Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Wed, 26 Aug 2026 02:31:04 +0000 Subject: [PATCH 22/23] style --- .../parquet/experimental/hybrid_scan_impl.cpp | 8 +++--- cpp/src/io/parquet/reader_impl.hpp | 24 ++++++++--------- cpp/src/io/parquet/reader_impl_helpers.cpp | 26 +++++++++---------- cpp/tests/io/parquet_reader_test.cpp | 2 +- 4 files changed, 30 insertions(+), 30 deletions(-) diff --git a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp index ba377c38cb75..781a53111776 100644 --- a/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp +++ b/cpp/src/io/parquet/experimental/hybrid_scan_impl.cpp @@ -179,10 +179,10 @@ hybrid_scan_reader_impl::hybrid_scan_reader_impl( cudf::host_span const> footer_bytes, parquet_reader_options const& options) { - _metadata = std::make_shared( - footer_bytes, - options.is_enabled_use_arrow_schema(), - options.is_enabled_allow_mismatched_pq_schemas()); + _metadata = + std::make_shared(footer_bytes, + options.is_enabled_use_arrow_schema(), + options.is_enabled_allow_mismatched_pq_schemas()); _extended_metadata = static_cast(_metadata.get()); } diff --git a/cpp/src/io/parquet/reader_impl.hpp b/cpp/src/io/parquet/reader_impl.hpp index 30d7550264f0..c8209c8b5bda 100644 --- a/cpp/src/io/parquet/reader_impl.hpp +++ b/cpp/src/io/parquet/reader_impl.hpp @@ -427,19 +427,19 @@ class reader_impl { not(has_cols_from_mismatched_sources(options) and _metadata->get_num_sources() > 1); } - private: + private: /** - * @brief Check if the user has specified columns from mismatched sources - * - * @param options Reader options - * @return True if the user has specified columns from mismatched sources - */ - [[nodiscard]] bool has_cols_from_mismatched_sources(parquet_reader_options const& options) const - { - return (options.get_column_names().has_value() or - options.get_column_field_ids().has_value()) and - options.is_enabled_allow_mismatched_pq_schemas(); - } + * @brief Check if the user has specified columns from mismatched sources + * + * @param options Reader options + * @return True if the user has specified columns from mismatched sources + */ + [[nodiscard]] bool has_cols_from_mismatched_sources(parquet_reader_options const& options) const + { + return (options.get_column_names().has_value() or + options.get_column_field_ids().has_value()) and + options.is_enabled_allow_mismatched_pq_schemas(); + } protected: /** diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index a302476f6ed9..bbf3b9288cb8 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -2103,19 +2103,19 @@ aggregate_reader_metadata::select_columns( }; // Compares two schema elements to be equal except their number of children - auto const equal_to_except_num_children = [selection_mode](SchemaElement const& lhs, - SchemaElement const& rhs) { - // Match by field ID if enabled, otherwise match by name - auto const match_schema_by_field_id = selection_mode == column_selection_mode::BY_FIELD_ID; - auto const names_match = - (match_schema_by_field_id and lhs.field_id.has_value() and rhs.field_id.has_value()) - ? lhs.field_id == rhs.field_id - : lhs.name == rhs.name; - return lhs.type == rhs.type and lhs.converted_type == rhs.converted_type and - lhs.type_length == rhs.type_length and names_match and - lhs.decimal_scale == rhs.decimal_scale and - lhs.decimal_precision == rhs.decimal_precision and lhs.field_id == rhs.field_id; - }; + auto const equal_to_except_num_children = + [selection_mode, case_sensitive_names](SchemaElement const& lhs, SchemaElement const& rhs) { + // Match by field ID if enabled, otherwise match by name + auto const match_schema_by_field_id = selection_mode == column_selection_mode::BY_FIELD_ID; + auto const names_match = + (match_schema_by_field_id and lhs.field_id.has_value() and rhs.field_id.has_value()) + ? lhs.field_id == rhs.field_id + : are_column_paths_equal(lhs.name, rhs.name, case_sensitive_names); + return lhs.type == rhs.type and lhs.converted_type == rhs.converted_type and + lhs.type_length == rhs.type_length and names_match and + lhs.decimal_scale == rhs.decimal_scale and + lhs.decimal_precision == rhs.decimal_precision and lhs.field_id == rhs.field_id; + }; // Maps a projected column's schema_idx in the zeroth per_file_metadata (source) to the // corresponding schema_idx in src_idx'th per_file_metadata (destination). The projected diff --git a/cpp/tests/io/parquet_reader_test.cpp b/cpp/tests/io/parquet_reader_test.cpp index 6f86cd4b5e69..23b8807d21ef 100644 --- a/cpp/tests/io/parquet_reader_test.cpp +++ b/cpp/tests/io/parquet_reader_test.cpp @@ -6325,4 +6325,4 @@ TEST_F(ParquetReaderTest, CaseInsensitiveMismatchedSchemasPropagateNullability) {1, 2, 3, 4, 5}, {true, true, true, true, false}}; auto const result = cudf::io::read_parquet(options); CUDF_TEST_EXPECT_COLUMNS_EQUAL(result.tbl->view().column(0), expected); -} \ No newline at end of file +} From f0bfb56fcbf1b17ec21e0c0872ebf0919ab86b0c Mon Sep 17 00:00:00 2001 From: Muhammad Haseeb Date: Fri, 28 Aug 2026 20:33:06 +0000 Subject: [PATCH 23/23] clang format for the billionth time --- cpp/src/io/parquet/reader_impl_helpers.cpp | 26 +++++++++++----------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/cpp/src/io/parquet/reader_impl_helpers.cpp b/cpp/src/io/parquet/reader_impl_helpers.cpp index c014d07fd5ac..ad5d1236aa63 100644 --- a/cpp/src/io/parquet/reader_impl_helpers.cpp +++ b/cpp/src/io/parquet/reader_impl_helpers.cpp @@ -2104,19 +2104,19 @@ aggregate_reader_metadata::select_columns( }; // Compares two schema elements to be equal except their number of children - auto const equal_to_except_num_children = - [selection_mode, case_sensitive_names](SchemaElement const& lhs, SchemaElement const& rhs) { - // Match by field ID if enabled, otherwise match by name - auto const match_schema_by_field_id = selection_mode == column_selection_mode::BY_FIELD_ID; - auto const names_match = - (match_schema_by_field_id and lhs.field_id.has_value() and rhs.field_id.has_value()) - ? lhs.field_id == rhs.field_id - : are_column_paths_equal(lhs.name, rhs.name, case_sensitive_names); - return lhs.type == rhs.type and lhs.converted_type == rhs.converted_type and - lhs.type_length == rhs.type_length and names_match and - lhs.decimal_scale == rhs.decimal_scale and - lhs.decimal_precision == rhs.decimal_precision and lhs.field_id == rhs.field_id; - }; + auto const equal_to_except_num_children = [selection_mode, case_sensitive_names]( + SchemaElement const& lhs, SchemaElement const& rhs) { + // Match by field ID if enabled, otherwise match by name + auto const match_schema_by_field_id = selection_mode == column_selection_mode::BY_FIELD_ID; + auto const names_match = + (match_schema_by_field_id and lhs.field_id.has_value() and rhs.field_id.has_value()) + ? lhs.field_id == rhs.field_id + : are_column_paths_equal(lhs.name, rhs.name, case_sensitive_names); + return lhs.type == rhs.type and lhs.converted_type == rhs.converted_type and + lhs.type_length == rhs.type_length and names_match and + lhs.decimal_scale == rhs.decimal_scale and + lhs.decimal_precision == rhs.decimal_precision and lhs.field_id == rhs.field_id; + }; // Maps a projected column's schema_idx in the zeroth per_file_metadata (source) to the // corresponding schema_idx in src_idx'th per_file_metadata (destination). The projected