diff --git a/cpp/src/io/parquet/io_utils/parquet_io_utils.cpp b/cpp/src/io/parquet/io_utils/parquet_io_utils.cpp index a28f00c9aca1..bd685d238a7e 100644 --- a/cpp/src/io/parquet/io_utils/parquet_io_utils.cpp +++ b/cpp/src/io/parquet/io_utils/parquet_io_utils.cpp @@ -323,8 +323,12 @@ fetch_byte_ranges_to_device_async_impl( auto const& byte_ranges = byte_ranges_per_source[source_idx]; // Total buffer size required for column chunks of this source + auto const source_size = datasources[source_idx].get().size(); auto const buffer_size = std::accumulate( - byte_ranges.begin(), byte_ranges.end(), std::size_t{0}, [](auto acc, auto const& range) { + byte_ranges.begin(), byte_ranges.end(), std::size_t{0}, [&](auto acc, auto const& range) { + CUDF_EXPECTS( + static_cast(range.offset()) + static_cast(range.size()) <= source_size, + "Byte range exceeds datasource size"); return acc + range.size(); }); diff --git a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp index eee71f349ec7..505c5da05a57 100644 --- a/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp +++ b/cpp/tests/io/experimental/hybrid_scan_filters_test.cpp @@ -1774,6 +1774,45 @@ TEST_F(HybridScanFiltersTest, RowGroupPasses) } } +TEST_F(HybridScanFiltersTest, FetchByteRangesInvalidRanges) +{ + std::vector data(1024); + auto const datasource = + cudf::io::datasource::create(cudf::host_span(data.data(), data.size())); + auto const stream = cudf::get_default_stream(); + auto const mr = cudf::get_current_device_resource_ref(); + + EXPECT_THROW( + cudf::io::parquet::fetch_byte_ranges_to_device_async( + *datasource, + std::vector{cudf::io::text::byte_range_info{-1, 16}}, + stream, + mr), + cudf::logic_error); + + EXPECT_THROW( + cudf::io::parquet::fetch_byte_ranges_to_device_async( + *datasource, + std::vector{cudf::io::text::byte_range_info{512, 1024}}, + stream, + mr), + cudf::logic_error); + + EXPECT_THROW( + cudf::io::parquet::fetch_byte_ranges_to_device_async( + *datasource, + std::vector{cudf::io::text::byte_range_info{0, -1}}, + stream, + mr), + cudf::logic_error); + + EXPECT_NO_THROW(cudf::io::parquet::fetch_byte_ranges_to_device_async( + *datasource, + std::vector{cudf::io::text::byte_range_info{1023, 1}}, + stream, + mr)); +} + class DictionaryFilterGapTest : public HybridScanFiltersTest, public ::testing::WithParamInterface {}; diff --git a/docs/cudf/source/pylibcudf/api_docs/io/experimental.rst b/docs/cudf/source/pylibcudf/api_docs/io/experimental.rst new file mode 100644 index 000000000000..92e9360fdb42 --- /dev/null +++ b/docs/cudf/source/pylibcudf/api_docs/io/experimental.rst @@ -0,0 +1,8 @@ +============ +Experimental +============ + +APIs in this namespace are experimental and may change without warning in the future. + +.. automodule:: pylibcudf.io.experimental + :members: diff --git a/docs/cudf/source/pylibcudf/api_docs/io/index.rst b/docs/cudf/source/pylibcudf/api_docs/io/index.rst index 15a87175325f..45b4def70051 100644 --- a/docs/cudf/source/pylibcudf/api_docs/io/index.rst +++ b/docs/cudf/source/pylibcudf/api_docs/io/index.rst @@ -17,9 +17,11 @@ I/O Functions avro csv + experimental json orc parquet + parquet_io_utils parquet_metadata text timezone diff --git a/docs/cudf/source/pylibcudf/api_docs/io/parquet_io_utils.rst b/docs/cudf/source/pylibcudf/api_docs/io/parquet_io_utils.rst new file mode 100644 index 000000000000..3f01b0493fd4 --- /dev/null +++ b/docs/cudf/source/pylibcudf/api_docs/io/parquet_io_utils.rst @@ -0,0 +1,6 @@ +================ +Parquet IO Utils +================ + +.. automodule:: pylibcudf.io.parquet_io_utils + :members: diff --git a/python/pylibcudf/pylibcudf/io/CMakeLists.txt b/python/pylibcudf/pylibcudf/io/CMakeLists.txt index 089ea8d0e8d9..65bb31d908ff 100644 --- a/python/pylibcudf/pylibcudf/io/CMakeLists.txt +++ b/python/pylibcudf/pylibcudf/io/CMakeLists.txt @@ -1,12 +1,12 @@ # ============================================================================= # cmake-format: off -# SPDX-FileCopyrightText: Copyright (c) 2024-2025, NVIDIA CORPORATION. +# SPDX-FileCopyrightText: Copyright (c) 2024-2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. # SPDX-License-Identifier: Apache-2.0 # cmake-format: on # ============================================================================= set(cython_sources avro.pyx csv.pyx datasource.pyx json.pyx orc.pyx parquet.pyx - parquet_metadata.pyx text.pyx timezone.pyx types.pyx + parquet_io_utils.pyx parquet_metadata.pyx text.pyx timezone.pyx types.pyx ) set(linked_libraries cudf::cudf) diff --git a/python/pylibcudf/pylibcudf/io/__init__.py b/python/pylibcudf/pylibcudf/io/__init__.py index a6a0ebad3a1e..1f0a0a218199 100644 --- a/python/pylibcudf/pylibcudf/io/__init__.py +++ b/python/pylibcudf/pylibcudf/io/__init__.py @@ -9,6 +9,7 @@ json, orc, parquet, + parquet_io_utils, parquet_metadata, text, timezone, @@ -30,6 +31,7 @@ "json", "orc", "parquet", + "parquet_io_utils", "parquet_metadata", "text", "timezone", diff --git a/python/pylibcudf/pylibcudf/io/parquet_io_utils.pxd b/python/pylibcudf/pylibcudf/io/parquet_io_utils.pxd new file mode 100644 index 000000000000..e8d872e08f61 --- /dev/null +++ b/python/pylibcudf/pylibcudf/io/parquet_io_utils.pxd @@ -0,0 +1,18 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from pylibcudf.io.text cimport ByteRangeInfo +from pylibcudf.io.types cimport SourceInfo +from rmm.pylibrmm.memory_resource cimport DeviceMemoryResource + +cpdef list fetch_byte_ranges_to_device( + SourceInfo source_info, + list byte_ranges, + object stream=*, + DeviceMemoryResource mr=*, +) + +cpdef bytes fetch_page_index_to_host( + SourceInfo source_info, + ByteRangeInfo page_index_range, +) diff --git a/python/pylibcudf/pylibcudf/io/parquet_io_utils.pyi b/python/pylibcudf/pylibcudf/io/parquet_io_utils.pyi new file mode 100644 index 000000000000..1a18ab72f5d5 --- /dev/null +++ b/python/pylibcudf/pylibcudf/io/parquet_io_utils.pyi @@ -0,0 +1,22 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from rmm.pylibrmm.memory_resource import DeviceMemoryResource + +from pylibcudf.gpumemoryview import gpumemoryview +from pylibcudf.io.text import ByteRangeInfo +from pylibcudf.io.types import SourceInfo +from pylibcudf.utils import CudaStreamLike + +__all__ = ["fetch_byte_ranges_to_device", "fetch_page_index_to_host"] + +def fetch_byte_ranges_to_device( + source_info: SourceInfo, + byte_ranges: list[ByteRangeInfo], + stream: CudaStreamLike | None = None, + mr: DeviceMemoryResource | None = None, +) -> list[gpumemoryview]: ... +def fetch_page_index_to_host( + source_info: SourceInfo, + page_index_range: ByteRangeInfo, +) -> bytes: ... diff --git a/python/pylibcudf/pylibcudf/io/parquet_io_utils.pyx b/python/pylibcudf/pylibcudf/io/parquet_io_utils.pyx new file mode 100644 index 000000000000..e54d879fa5c2 --- /dev/null +++ b/python/pylibcudf/pylibcudf/io/parquet_io_utils.pyx @@ -0,0 +1,157 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +"""IO utilities for Parquet.""" + +from libc.stddef cimport size_t +from libc.stdint cimport uint8_t, uintptr_t +from libcpp.memory cimport make_unique, unique_ptr +from libcpp.pair cimport pair +from libcpp.utility cimport move +from libcpp.vector cimport vector +from cython.operator cimport dereference + +from rmm.librmm.device_buffer cimport device_buffer +from rmm.pylibrmm.device_buffer cimport DeviceBuffer +from rmm.pylibrmm.memory_resource cimport DeviceMemoryResource +from rmm.pylibrmm.stream cimport Stream + +from pylibcudf.gpumemoryview cimport gpumemoryview +from pylibcudf.io.text cimport ByteRangeInfo +from pylibcudf.io.types cimport SourceInfo +from pylibcudf.libcudf.io.datasource cimport datasource, make_datasources +from pylibcudf.libcudf.io.parquet_io_utils cimport ( + const_byte_range_info, + const_uint8_t, + cpp_fetch_byte_ranges_to_device, + fetch_page_index_to_host as cpp_fetch_page_index_to_host, +) + +from pylibcudf.libcudf.io.text cimport byte_range_info +from pylibcudf.libcudf.utilities.span cimport device_span, host_span +from pylibcudf.utils cimport _get_memory_resource, _get_stream + +__all__ = ["fetch_byte_ranges_to_device", "fetch_page_index_to_host"] + + +cpdef list fetch_byte_ranges_to_device( + SourceInfo source_info, + list byte_ranges, + object stream=None, + DeviceMemoryResource mr=None, +): + """Fetch byte ranges from a Parquet source into device memory. + + Parameters + ---------- + source_info : SourceInfo + Source describing a single Parquet file. + byte_ranges : list[ByteRangeInfo] + Byte ranges to fetch, as returned by + :meth:`~pylibcudf.io.experimental.HybridScanReader.filter_column_chunks_byte_ranges`, + :meth:`~pylibcudf.io.experimental.HybridScanReader.payload_column_chunks_byte_ranges`, + or + :meth:`~pylibcudf.io.experimental.HybridScanReader.all_column_chunks_byte_ranges`. + stream : Stream, optional + CUDA stream. + mr : DeviceMemoryResource, optional + Device memory resource. + + Returns + ------- + list[gpumemoryview] + One view per byte range. Each view holds a reference to the + :class:`~rmm.DeviceBuffer` that owns its memory, keeping the + allocation alive for as long as the view is referenced. + + Raises + ------ + ValueError + If ``source_info`` does not describe exactly one source. + """ + cdef Stream _stream = _get_stream(stream) + cdef DeviceMemoryResource _mr = _get_memory_resource(mr) + cdef vector[unique_ptr[datasource]] sources = make_datasources(source_info.c_obj) + if sources.size() != 1: + raise ValueError( + f"fetch_byte_ranges_to_device requires exactly one source, " + f"got {sources.size()}" + ) + + cdef vector[byte_range_info] ranges_vec + cdef ByteRangeInfo bri + for bri in byte_ranges: + ranges_vec.push_back(bri.c_obj) + + cdef pair[vector[device_buffer], vector[device_span[const_uint8_t]]] fetched + with nogil: + fetched = cpp_fetch_byte_ranges_to_device( + dereference(sources[0]), + host_span[const_byte_range_info](ranges_vec.data(), ranges_vec.size()), + _stream.view(), + _mr.get_mr(), + ) + + if fetched.first.size() != 1: + raise RuntimeError( + f"Expected exactly one device buffer, got {fetched.first.size()}" + ) + cdef DeviceBuffer owner = DeviceBuffer.c_from_unique_ptr( + make_unique[device_buffer](move(fetched.first[0])), + _stream, + _mr, + ) + cdef gpumemoryview owner_gv = gpumemoryview(owner) + cdef uintptr_t base = owner_gv.ptr + cdef uintptr_t ptr + cdef size_t n + result = [] + for i in range(fetched.second.size()): + ptr = fetched.second[i].data() + n = fetched.second[i].size() + result.append(owner_gv.byte_slice(slice(ptr - base, ptr - base + n))) + return result + + +cpdef bytes fetch_page_index_to_host( + SourceInfo source_info, + ByteRangeInfo page_index_range, +): + """Fetch parquet page index bytes to host memory. + + Parameters + ---------- + source_info : SourceInfo + Source describing a single Parquet file. + page_index_range : ByteRangeInfo + Byte range of the page index, as returned by + :meth:`~pylibcudf.io.experimental.HybridScanReader.page_index_byte_range`. + + Returns + ------- + bytes + Raw page index bytes copied to Python host memory. + + Raises + ------ + ValueError + If ``source_info`` does not describe exactly one source. + """ + cdef vector[unique_ptr[datasource]] sources = make_datasources(source_info.c_obj) + if sources.size() != 1: + raise ValueError( + f"fetch_page_index_to_host requires exactly one source, " + f"got {sources.size()}" + ) + + cdef unique_ptr[datasource.buffer] buf + with nogil: + buf = move(cpp_fetch_page_index_to_host( + dereference(sources[0]), + (page_index_range).c_obj, + )) + + if buf.get() is NULL: + raise RuntimeError("fetch_page_index_to_host returned no buffer") + cdef const uint8_t* ptr = buf.get().data() + cdef size_t n = buf.get().size() + return bytes(ptr[:n]) diff --git a/python/pylibcudf/pylibcudf/io/parquet_metadata.pyx b/python/pylibcudf/pylibcudf/io/parquet_metadata.pyx index e6015786173b..f5bb4782a87a 100644 --- a/python/pylibcudf/pylibcudf/io/parquet_metadata.pyx +++ b/python/pylibcudf/pylibcudf/io/parquet_metadata.pyx @@ -465,7 +465,7 @@ cdef class FileMetaData: See Also -------- - read_parquet_footers + pylibcudf.io.parquet_metadata.read_parquet_footers Read one ``FileMetaData`` per source directly from :class:`pylibcudf.io.types.SourceInfo`. """ diff --git a/python/pylibcudf/pylibcudf/libcudf/io/datasource.pxd b/python/pylibcudf/pylibcudf/libcudf/io/datasource.pxd index 36c5bf928850..0e047b2fd836 100644 --- a/python/pylibcudf/pylibcudf/libcudf/io/datasource.pxd +++ b/python/pylibcudf/pylibcudf/libcudf/io/datasource.pxd @@ -1,6 +1,7 @@ -# SPDX-FileCopyrightText: Copyright (c) 2023-2026, NVIDIA CORPORATION. +# SPDX-FileCopyrightText: Copyright (c) 2023-2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. # SPDX-License-Identifier: Apache-2.0 from libc.stddef cimport size_t +from libc.stdint cimport uint8_t from libcpp.memory cimport unique_ptr from libcpp.vector cimport vector from pylibcudf.libcudf.io.types cimport source_info @@ -11,7 +12,9 @@ cdef extern from "cudf/io/datasource.hpp" \ namespace "cudf::io" nogil: cdef cppclass datasource: - pass + cppclass buffer: + const uint8_t* data() except +libcudf_exception_handler const + size_t size() except +libcudf_exception_handler const cdef vector[unique_ptr[datasource]] make_datasources( source_info info diff --git a/python/pylibcudf/pylibcudf/libcudf/io/parquet_io_utils.pxd b/python/pylibcudf/pylibcudf/libcudf/io/parquet_io_utils.pxd new file mode 100644 index 000000000000..47f871439b99 --- /dev/null +++ b/python/pylibcudf/pylibcudf/libcudf/io/parquet_io_utils.pxd @@ -0,0 +1,57 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from libc.stdint cimport uint8_t +from libcpp.memory cimport unique_ptr +from libcpp.pair cimport pair +from libcpp.vector cimport vector + +from rmm.librmm.cuda_stream_view cimport cuda_stream_view +from rmm.librmm.device_buffer cimport device_buffer +from rmm.librmm.memory_resource cimport device_async_resource_ref + +from pylibcudf.exception_handler cimport libcudf_exception_handler +from pylibcudf.libcudf.io.datasource cimport datasource +from pylibcudf.libcudf.io.text cimport byte_range_info +from pylibcudf.libcudf.utilities.span cimport device_span, host_span + +ctypedef const uint8_t const_uint8_t +ctypedef const byte_range_info const_byte_range_info + +cdef extern from * nogil: + """ + #include + #include + + static std::pair, + std::vector>> + cpp_fetch_byte_ranges_to_device( + cudf::io::datasource& datasource, + cudf::host_span byte_ranges, + rmm::cuda_stream_view stream, + rmm::device_async_resource_ref mr) + { + auto [buffers, spans, fut] = + cudf::io::parquet::fetch_byte_ranges_to_device_async( + datasource, byte_ranges, stream, mr); + // Block until the async fetch completes so the returned buffers/spans + // are fully populated. This also avoids exposing std::future to Cython. + fut.get(); + return {std::move(buffers), std::move(spans)}; + } + """ + pair[vector[device_buffer], vector[device_span[const_uint8_t]]] \ + cpp_fetch_byte_ranges_to_device( + datasource& source, + host_span[const_byte_range_info] byte_ranges, + cuda_stream_view stream, + device_async_resource_ref mr, + ) except +libcudf_exception_handler + +cdef extern from "cudf/io/parquet_io_utils.hpp" \ + namespace "cudf::io::parquet" nogil: + + unique_ptr[datasource.buffer] fetch_page_index_to_host( + datasource& ds, + byte_range_info page_index_bytes, + ) except +libcudf_exception_handler diff --git a/python/pylibcudf/pylibcudf/libcudf/utilities/span.pxd b/python/pylibcudf/pylibcudf/libcudf/utilities/span.pxd index f2bf388e4d4c..4f4669d7f6f8 100644 --- a/python/pylibcudf/pylibcudf/libcudf/utilities/span.pxd +++ b/python/pylibcudf/pylibcudf/libcudf/utilities/span.pxd @@ -1,5 +1,6 @@ -# SPDX-FileCopyrightText: Copyright (c) 2021-2025, NVIDIA CORPORATION. +# SPDX-FileCopyrightText: Copyright (c) 2021-2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. # SPDX-License-Identifier: Apache-2.0 +from libc.stddef cimport size_t from libcpp.vector cimport vector from pylibcudf.exception_handler cimport libcudf_exception_handler from pylibcudf.libcudf.types cimport size_type @@ -13,5 +14,6 @@ cdef extern from "cudf/utilities/span.hpp" namespace "cudf" nogil: cdef cppclass device_span[T]: device_span() noexcept - device_span(T *data, size_type size) noexcept - T *data() noexcept + device_span(T *data, size_t size) noexcept + T *data() noexcept const + size_t size() noexcept const