diff --git a/cpp/include/cudf/io/experimental/deletion_vectors.hpp b/cpp/include/cudf/io/experimental/deletion_vectors.hpp index af3026a671d8..fbab01669b67 100644 --- a/cpp/include/cudf/io/experimental/deletion_vectors.hpp +++ b/cpp/include/cudf/io/experimental/deletion_vectors.hpp @@ -48,6 +48,9 @@ struct deletion_vector_info { std::vector row_group_offsets; /// Number of rows in each row group to be read from the Parquet source(s) std::vector row_group_num_rows; + + /// Whether the roaring bitmaps represent retention vectors + bool are_retention_vectors = false; }; /** @@ -147,6 +150,7 @@ class chunked_parquet_reader { std::queue _deletion_vector_row_counts; size_t _start_row; bool _is_unspecified_row_group_data; + bool _are_retentions; rmm::cuda_stream_view _stream; rmm::device_async_resource_ref _mr; rmm::device_async_resource_ref _table_mr; diff --git a/cpp/src/io/parquet/experimental/deletion_vectors.cu b/cpp/src/io/parquet/experimental/deletion_vectors.cu index 1d5ada49c69a..7af0cd3b3cf5 100644 --- a/cpp/src/io/parquet/experimental/deletion_vectors.cu +++ b/cpp/src/io/parquet/experimental/deletion_vectors.cu @@ -104,11 +104,15 @@ namespace detail { deletion_vector_row_counts, stream, cudf::get_current_device_resource_ref()); + + auto const mask_type = deletion_vector_info.are_retention_vectors + ? cudf::detail::mask_type::RETENTION + : cudf::detail::mask_type::DELETION; + // Filter the table using the deletion vector - return table_with_metadata{ + return cudf::io::table_with_metadata{ // Supply user-provided mr to apply deletion mask to allocate output table's memory - cudf::detail::apply_mask( - table_with_index->view(), row_mask->view(), cudf::detail::mask_type::DELETION, stream, mr), + cudf::detail::apply_mask(table_with_index->view(), row_mask->view(), mask_type, stream, mr), std::move(metadata)}; } @@ -161,7 +165,7 @@ namespace detail { dv_row_counts_queue.push(deletion_vector_row_counts[i]); } - size_t deleted_rows = 0; + size_t matched_rows = 0; size_t remaining_rows = num_rows; size_t start_row = 0; @@ -176,14 +180,15 @@ namespace detail { is_row_group_data_unspecified, stream, cudf::get_current_device_resource_ref()); - deleted_rows += compute_partial_deleted_row_count( + matched_rows += compute_partial_deleted_row_count( row_index_column->view(), dv_queue, dv_row_counts_queue, stream); start_row += chunk_rows; remaining_rows -= chunk_rows; } - return deleted_rows; + // Bitmap hits are deleted rows for deletion vectors, retained rows for retention vectors + return deletion_vector_info.are_retention_vectors ? num_rows - matched_rows : matched_rows; } } // namespace detail @@ -200,6 +205,7 @@ chunked_parquet_reader::chunked_parquet_reader(std::size_t chunk_read_limit, rmm::device_async_resource_ref mr) : _start_row{0}, _is_unspecified_row_group_data{deletion_vector_info.row_group_offsets.empty()}, + _are_retentions{deletion_vector_info.are_retention_vectors}, _stream{stream}, _mr{mr}, // Use default mr for the internal chunked reader and row index column if we will @@ -317,10 +323,11 @@ table_with_metadata chunked_parquet_reader::read_chunk() _deletion_vector_row_counts, _stream, cudf::get_current_device_resource_ref()); - return table_with_metadata{ + auto const mask_type = + _are_retentions ? cudf::detail::mask_type::RETENTION : cudf::detail::mask_type::DELETION; + return cudf::io::table_with_metadata{ // Supply user-provided mr to apply deletion mask to allocate output table's memory - cudf::detail::apply_mask( - table_with_index->view(), row_mask->view(), cudf::detail::mask_type::DELETION, _stream, _mr), + cudf::detail::apply_mask(table_with_index->view(), row_mask->view(), mask_type, _stream, _mr), std::move(metadata)}; } diff --git a/cpp/tests/io/parquet_deletion_vectors_test.cpp b/cpp/tests/io/parquet_deletion_vectors_test.cpp index a41652105e52..c261fe172be7 100644 --- a/cpp/tests/io/parquet_deletion_vectors_test.cpp +++ b/cpp/tests/io/parquet_deletion_vectors_test.cpp @@ -1,5 +1,5 @@ /* - * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION. + * SPDX-FileCopyrightText: Copyright (c) 2025-2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. * SPDX-License-Identifier: Apache-2.0 */ @@ -136,14 +136,17 @@ auto build_expected_row_indices(cudf::host_span row_group_off * @param num_rows Number of rows in the table * @param deletion_probability The probability of a row being deleted * @param row_indices Host vector of row indices + * @param are_retention_vectors Whether to add retained, rather than deleted, row indices to the + * bitmap * - * @return A pair of a deletion vector and a host row mask vector + * @return A pair of a roaring bitmap and a host row mask vector */ -auto build_deletion_vector_and_expected_row_mask(cudf::size_type num_rows, - float deletion_probability, - cudf::host_span row_indices, - rmm::cuda_stream_view stream, - rmm::device_async_resource_ref mr) +auto build_roaring_bitmap_and_expected_row_mask(cudf::size_type num_rows, + float deletion_probability, + cudf::host_span row_indices, + rmm::cuda_stream_view stream, + rmm::device_async_resource_ref mr, + bool are_retention_vectors = false) { static constexpr auto seed = 0xbaLL; std::mt19937 engine{seed}; @@ -164,8 +167,9 @@ auto build_deletion_vector_and_expected_row_mask(cudf::size_type num_rows, std::for_each(cuda::counting_iterator(0), cuda::counting_iterator(num_rows), [&](auto row_idx) { - // Insert provided host row index if the row is deleted in the row mask - if (not expected_row_mask[row_idx]) { + // For retention vectors, retain rows selected in the row mask. For deletion + // vectors, delete the remaining rows. + if (expected_row_mask[row_idx] == are_retention_vectors) { roaring::api::roaring64_bitmap_add_bulk( deletion_vector, &roaring64_context, row_indices[row_idx]); } @@ -222,7 +226,7 @@ std::unique_ptr build_expected_table( * well as the `cudf::io::parquet::experimental::chunked_parquet_reader` */ template -void test_read_parquet_and_apply_deletion_vector( +void test_read_parquet_and_apply_mask( cudf::host_span parquet_buffer, cudf::io::parquet::experimental::deletion_vector_info const& deletion_vector_info, cudf::table_view const& input_table_view, @@ -245,6 +249,7 @@ void test_read_parquet_and_apply_deletion_vector( .deletion_vector_row_counts = std::vector(num_concat, input_table_view.num_rows()), .row_group_offsets = {}, .row_group_num_rows = {}, + .are_retention_vectors = deletion_vector_info.are_retention_vectors, }; // Vector to hold the Parquet buffer spans @@ -296,8 +301,13 @@ void test_read_parquet_and_apply_deletion_vector( local_expected_row_indices, cudf::type_id::UINT64, stream, mr); auto [local_deletion_vector, local_expected_row_mask_column] = - build_deletion_vector_and_expected_row_mask( - num_input_rows, deletion_probability, local_expected_row_indices, stream, mr); + build_roaring_bitmap_and_expected_row_mask( + num_input_rows, + deletion_probability, + local_expected_row_indices, + stream, + mr, + final_deletion_vector_info.are_retention_vectors); // Insert the expected table, the corresponding deletion vector and its data span tables.emplace_back(build_expected_table(input_table_view, @@ -383,37 +393,58 @@ TEST_F(ParquetDeletionVectorsTest, NoRowIndexColumn) auto expected_row_index_column = build_column_from_host_data( expected_row_indices, cudf::type_id::UINT64, stream, mr); - // Build deletion vector and the expected row mask column - auto [deletion_vector, expected_row_mask_column] = build_deletion_vector_and_expected_row_mask( - num_rows, deletion_probability, expected_row_indices, stream, mr); - // Use num_concat = 1 here since the row index column is simply a sequence and input table // concatenation won't properly reset it. auto constexpr num_concat = 1; - cudf::io::parquet::experimental::deletion_vector_info deletion_vector_info{ - .serialized_roaring_bitmaps = {deletion_vector}, - .deletion_vector_row_counts = {input_table->view().num_rows()}}; - test_read_parquet_and_apply_deletion_vector(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); + + // Build deletion vector and the expected row mask column + { + auto [deletion_vector, expected_row_mask_column] = build_roaring_bitmap_and_expected_row_mask( + num_rows, deletion_probability, expected_row_indices, stream, mr); + auto deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{ + .serialized_roaring_bitmaps = {deletion_vector}, + .deletion_vector_row_counts = {input_table->view().num_rows()}}; + test_read_parquet_and_apply_mask(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); + } + + // Retention vectors retain the rows present in the bitmap. + { + auto constexpr are_retention_vectors = true; + auto [retention_vector, retention_row_mask_column] = build_roaring_bitmap_and_expected_row_mask( + num_rows, deletion_probability, expected_row_indices, stream, mr, are_retention_vectors); + auto retention_vector_info = cudf::io::parquet::experimental::deletion_vector_info{ + .serialized_roaring_bitmaps = {retention_vector}, + .deletion_vector_row_counts = {input_table->view().num_rows()}, + .are_retention_vectors = are_retention_vectors}; + test_read_parquet_and_apply_mask(parquet_buffer, + retention_vector_info, + input_table->view(), + retention_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); + } + // Test no row index column and no deletion vector { // Build expected row mask column containing all true values auto expected_row_mask = thrust::host_vector(num_rows, true); auto expected_row_mask_column = build_column_from_host_data(expected_row_mask, cudf::type_id::BOOL8, stream, mr); - deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{}; - test_read_parquet_and_apply_deletion_vector(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); + auto deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{}; + test_read_parquet_and_apply_mask(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); } } @@ -477,7 +508,7 @@ TEST_F(ParquetDeletionVectorsTest, CustomRowIndexColumn) expected_row_indices, cudf::type_id::UINT64, stream, mr); // Build deletion vector and the expected row mask column - auto [deletion_vector, expected_row_mask_column] = build_deletion_vector_and_expected_row_mask( + auto [deletion_vector, expected_row_mask_column] = build_roaring_bitmap_and_expected_row_mask( num_rows, deletion_probability, expected_row_indices, stream, mr); // Don't concatenate the input table and test with single deletion vector @@ -486,38 +517,58 @@ TEST_F(ParquetDeletionVectorsTest, CustomRowIndexColumn) .deletion_vector_row_counts = {input_table->view().num_rows()}, .row_group_offsets = row_group_offsets, .row_group_num_rows = row_group_num_rows}; - test_read_parquet_and_apply_deletion_vector<1>(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); + test_read_parquet_and_apply_mask<1>(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); // Concatenate input table and test with multiple deletion vectors - test_read_parquet_and_apply_deletion_vector<4>(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); + test_read_parquet_and_apply_mask<4>(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); // Concatenate input table and test with many deletion vectors (>= stream fork threshold of 8) - test_read_parquet_and_apply_deletion_vector<8>(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); - test_read_parquet_and_apply_deletion_vector<16>(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); + test_read_parquet_and_apply_mask<8>(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); + test_read_parquet_and_apply_mask<16>(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); + + // Retention vectors with custom row indices. + { + auto constexpr are_retention_vectors = true; + auto [retention_vector, retention_row_mask_column] = build_roaring_bitmap_and_expected_row_mask( + num_rows, deletion_probability, expected_row_indices, stream, mr, are_retention_vectors); + deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{ + .serialized_roaring_bitmaps = {retention_vector}, + .deletion_vector_row_counts = {input_table->view().num_rows()}, + .row_group_offsets = row_group_offsets, + .row_group_num_rows = row_group_num_rows, + .are_retention_vectors = are_retention_vectors}; + test_read_parquet_and_apply_mask<1>(parquet_buffer, + deletion_vector_info, + input_table->view(), + retention_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); + } // Test custom row index column and no deletion vector { @@ -528,13 +579,13 @@ TEST_F(ParquetDeletionVectorsTest, CustomRowIndexColumn) auto deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{ .row_group_offsets = row_group_offsets, .row_group_num_rows = row_group_num_rows}; - test_read_parquet_and_apply_deletion_vector<1>(parquet_buffer, - deletion_vector_info, - input_table->view(), - expected_row_mask_column->view(), - expected_row_index_column->view(), - stream, - mr); + test_read_parquet_and_apply_mask<1>(parquet_buffer, + deletion_vector_info, + input_table->view(), + expected_row_mask_column->view(), + expected_row_index_column->view(), + stream, + mr); } } @@ -562,24 +613,28 @@ TEST_F(DeletionVectorsCountTests, NoRowIndex) auto row_indices = thrust::host_vector(num_rows); std::iota(row_indices.begin(), row_indices.end(), size_t{0}); - auto [deletion_vector, expected_row_mask_column] = build_deletion_vector_and_expected_row_mask( - num_rows, deletion_probability, row_indices, stream, mr); + for (auto const are_retention_vectors : {false, true}) { + auto [deletion_vector, expected_row_mask_column] = build_roaring_bitmap_and_expected_row_mask( + num_rows, deletion_probability, row_indices, stream, mr, are_retention_vectors); - auto const expected_row_mask = cudf::detail::make_host_vector( - cudf::device_span(expected_row_mask_column->view().data(), num_rows), stream); - auto const expected_deleted = - std::count(expected_row_mask.begin(), expected_row_mask.end(), false); + auto const expected_row_mask = cudf::detail::make_host_vector( + cudf::device_span(expected_row_mask_column->view().data(), num_rows), + stream); + auto const expected_deleted = + std::count(expected_row_mask.begin(), expected_row_mask.end(), false); - auto deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{ - .serialized_roaring_bitmaps = {deletion_vector}, - .deletion_vector_row_counts = {num_rows}, - .row_group_offsets = {}, - .row_group_num_rows = {}}; - - for (auto chunk_size : {num_rows, num_rows / 2}) { - auto const result = cudf::io::parquet::experimental::compute_num_deleted_rows( - deletion_vector_info, chunk_size, stream); - EXPECT_EQ(result, expected_deleted); + auto deletion_vector_info = cudf::io::parquet::experimental::deletion_vector_info{ + .serialized_roaring_bitmaps = {deletion_vector}, + .deletion_vector_row_counts = {num_rows}, + .row_group_offsets = {}, + .row_group_num_rows = {}, + .are_retention_vectors = are_retention_vectors}; + + for (auto chunk_size : {num_rows, num_rows / 2}) { + auto const result = cudf::io::parquet::experimental::compute_num_deleted_rows( + deletion_vector_info, chunk_size, stream); + EXPECT_EQ(result, expected_deleted); + } } } @@ -628,7 +683,7 @@ TEST_F(DeletionVectorsCountTests, CustomRowIndex) auto expected_row_indices = build_expected_row_indices(row_group_offsets, row_group_num_rows, num_rows); - auto [deletion_vector, expected_row_mask_column] = build_deletion_vector_and_expected_row_mask( + auto [deletion_vector, expected_row_mask_column] = build_roaring_bitmap_and_expected_row_mask( num_rows, deletion_probability, expected_row_indices, stream, mr); auto const expected_row_mask = cudf::detail::make_host_vector( @@ -681,7 +736,7 @@ TEST_F(DeletionVectorsCountTests, MultipleDeletionVectors) auto local_indices = cudf::host_span(expected_row_indices.data() + span_start, num_rows_per_dv); - auto [dv, mask_col] = build_deletion_vector_and_expected_row_mask( + auto [dv, mask_col] = build_roaring_bitmap_and_expected_row_mask( num_rows_per_dv, deletion_probability, local_indices, stream, mr); auto const host_mask = cudf::detail::make_host_vector(