diff --git a/conda/recipes/libcudf/recipe.yaml b/conda/recipes/libcudf/recipe.yaml index 6dd71caec399..84455e556a93 100644 --- a/conda/recipes/libcudf/recipe.yaml +++ b/conda/recipes/libcudf/recipe.yaml @@ -319,6 +319,8 @@ outputs: - cuda-version =${{ cuda_version }} - cuda-nvtx-dev - cuda-cudart-dev + - cuda-nvrtc-dev + - libnvjitlink-dev run: - ${{ pin_subpackage("libcudf", exact=True) }} - ${{ pin_compatible("cuda-version", upper_bound="x", lower_bound="x") }} diff --git a/cpp/CMakeLists.txt b/cpp/CMakeLists.txt index 26b7ee68aad7..89e54b2ae11d 100644 --- a/cpp/CMakeLists.txt +++ b/cpp/CMakeLists.txt @@ -198,10 +198,6 @@ set_property( ) message(VERBOSE "CUDF: LIBCUDF_LOGGING_LEVEL = '${LIBCUDF_LOGGING_LEVEL}'.") -if(NOT CUDF_GENERATED_INCLUDE_DIR) - set(CUDF_GENERATED_INCLUDE_DIR ${CUDF_BINARY_DIR}) -endif() - # ################################################################################################## # * linter configuration --------------------------------------------------------------------------- if(CUDF_CLANG_TIDY) @@ -431,6 +427,24 @@ if(NOT BUILD_SHARED_LIBS) endif() endif() +set(LIBCUDF_FRAGMENT_LINK_LIBRARIES + CCCL::CCCL + rapids_logger::rapids_logger + rmm::rmm + $ + $ + $ + ZLIB::ZLIB + nvcomp::nvcomp + kvikio::kvikio + nanoarrow::nanoarrow + zstd +) +set(LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES "$" + "$" +) +set(LIBCUDF_FRAGMENT_COMPILE_OPTIONS "$<$:${CUDF_CUDA_FLAGS}>") + rtcx_add_embed(cudf_cuda_embed) rtcx_embed_includes( @@ -479,50 +493,19 @@ foreach(INC_DIR IN LISTS LIBCUDACXX_RAW_INCLUDE_DIRS) ) endforeach() -rtcx_embed( - cudf_cuda_embed COMPRESSION zstd OUTPUT_DIRECTORY "${CUDF_GENERATED_INCLUDE_DIR}/rtcx_embed" -) +rtcx_embed(cudf_cuda_embed COMPRESSION zstd OUTPUT_DIRECTORY "${CMAKE_CURRENT_BINARY_DIR}/embed") rtcx_add_embed(cudf_fragments) -list(APPEND CUDF_PRECOMPILE_PHYSICAL_TYPES uint8_t uint16_t uint32_t uint64_t numeric::decimal32 - numeric::decimal64 numeric::decimal128 -) - -foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) - set(FRAGMENT_NAME transform_kernel) - get_property( - FILE_INDEX - TARGET cudf_fragments__embed_props - PROPERTY EMBED_FILE_INDEX - ) - set(VARIANT_NAME transform_kernel_${FILE_INDEX}) - set(INSTANCE - "cudf::jit::transform_kernel>, cudf::jit::type_list>>" +# This function precompiles the unary and binary transform kernel fragments for all fixed-width +# types. This helps amortize the cost of JIT compilation for kernels that have matching signatures +# and reduces the amount of time spent on source-based runtime JIT compilation for LTO-based JIT. +function(precompile_fixed_width_kernel_fragments) + list(APPEND CUDF_PRECOMPILE_PHYSICAL_TYPES uint8_t uint16_t uint32_t uint64_t numeric::decimal32 + numeric::decimal64 numeric::decimal128 ) - add_fragment( - cudf_fragments - FRAGMENT - ${VARIANT_NAME} - SOURCE - src/transform/jit/kernel.cu - KERNEL_INSTANCE - ${INSTANCE} - UDF_TYPE - "int(${TYPE} *, ${TYPE})" - DEFINITIONS - CUDF_LTO_MODE - ARRAY_IDS - ${FRAGMENT_NAME}_FILE_INDEX - ${FRAGMENT_NAME}_INSTANCE - ARRAY_VALUES - ${FILE_INDEX} - "${INSTANCE}" - ) -endforeach() - -foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) - foreach(RHS_IS_SCALAR IN ITEMS "false" "true") + # Pre-compile unary-op kernel fragments for fixed-width types. + foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) set(FRAGMENT_NAME transform_kernel) get_property( FILE_INDEX @@ -531,7 +514,7 @@ foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) ) set(VARIANT_NAME transform_kernel_${FILE_INDEX}) set(INSTANCE - "cudf::jit::transform_kernel, cudf::jit::column_accessor<1ULL, cudf::column_device_view_core, ${TYPE}, ${RHS_IS_SCALAR}, 0>>, cudf::jit::type_list>>" + "cudf::jit::transform_kernel>, cudf::jit::type_list>>" ) add_fragment( cudf_fragments @@ -542,7 +525,7 @@ foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) KERNEL_INSTANCE ${INSTANCE} UDF_TYPE - "int(${TYPE} *, ${TYPE}, ${TYPE})" + "int(${TYPE} *, ${TYPE})" DEFINITIONS CUDF_LTO_MODE ARRAY_IDS @@ -551,13 +534,137 @@ foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) ARRAY_VALUES ${FILE_INDEX} "${INSTANCE}" + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} ) + endforeach() -endforeach() -rtcx_embed( - cudf_fragments COMPRESSION none OUTPUT_DIRECTORY "${CUDF_GENERATED_INCLUDE_DIR}/rtcx_embed" -) + # Precompile binary-op kernel fragments for fixed-width types. To minimize binary size impact, + # only the RHS-scalar variants are precompiled. + foreach(TYPE IN ITEMS ${CUDF_PRECOMPILE_PHYSICAL_TYPES}) + foreach(RHS_IS_SCALAR IN ITEMS "false" "true") + set(FRAGMENT_NAME transform_kernel) + get_property( + FILE_INDEX + TARGET cudf_fragments__embed_props + PROPERTY EMBED_FILE_INDEX + ) + set(VARIANT_NAME transform_kernel_${FILE_INDEX}) + set(INSTANCE + "cudf::jit::transform_kernel, cudf::jit::column_accessor<1ULL, cudf::column_device_view_core, ${TYPE}, ${RHS_IS_SCALAR}, 0>>, cudf::jit::type_list>>" + ) + add_fragment( + cudf_fragments + FRAGMENT + ${VARIANT_NAME} + SOURCE + src/transform/jit/kernel.cu + KERNEL_INSTANCE + ${INSTANCE} + UDF_TYPE + "int(${TYPE} *, ${TYPE}, ${TYPE})" + DEFINITIONS + CUDF_LTO_MODE + ARRAY_IDS + ${FRAGMENT_NAME}_FILE_INDEX + ${FRAGMENT_NAME}_INSTANCE + ARRAY_VALUES + ${FILE_INDEX} + "${INSTANCE}" + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + ) + endforeach() + endforeach() +endfunction() + +# This function precompiles kernel fragments for string transforms. It precompiles kernel fragments +# for common string transform operations: multi-extract, multi-split, string-parsing. This helps +# amortize the cost of JIT compilation for kernels that have matching signatures +function(precompile_string_kernel_fragments) + list(APPEND CUDF_PRECOMPILE_STRING_OUTPUT_PHYSICAL_TYPES uint8_t uint16_t uint32_t uint64_t + cuda::std::span + ) + list( + APPEND + CUDF_PRECOMPILE_STRING_OUTPUT_COLUMN_TYPES + cudf::mutable_column_device_view_core + cudf::mutable_column_device_view_core + cudf::mutable_column_device_view_core + cudf::mutable_column_device_view_core + cudf::jit::mutable_strings_column_device_view + ) + foreach(OUTPUT_TYPE OUTPUT_COLUMN_TYPE IN ZIP_LISTS CUDF_PRECOMPILE_STRING_OUTPUT_PHYSICAL_TYPES + CUDF_PRECOMPILE_STRING_OUTPUT_COLUMN_TYPES + ) + # Pre-compile for up to 8 output columns + foreach(OUTPUT_INDICES IN ITEMS "0" "0;1" "0;1;2" "0;1;2;3" "0;1;2;3;4" "0;1;2;3;4;5" + "0;1;2;3;4;5;6" "0;1;2;3;4;5;6;7" + ) + set(FRAGMENT_NAME transform_kernel) + get_property( + FILE_INDEX + TARGET cudf_fragments__embed_props + PROPERTY EMBED_FILE_INDEX + ) + set(VARIANT_NAME transform_kernel_${FILE_INDEX}) + set(OUTPUT_ACCESSORS "") + set(OUTPUT_POINTERS "") + foreach(OUTPUT_INDEX IN LISTS OUTPUT_INDICES) + list( + APPEND + OUTPUT_ACCESSORS + "cudf::jit::column_accessor<${OUTPUT_INDEX}ULL, ${OUTPUT_COLUMN_TYPE}, ${OUTPUT_TYPE}, false, 0>" + ) + list(APPEND OUTPUT_POINTERS "${OUTPUT_TYPE} *") + endforeach() + list(JOIN OUTPUT_ACCESSORS " ," OUTPUT_ACCESSORS_STR) + list(JOIN OUTPUT_POINTERS " ," OUTPUT_POINTERS_STR) + set(INSTANCE + "cudf::jit::transform_kernel>, cudf::jit::type_list<${OUTPUT_ACCESSORS_STR}>>" + ) + add_fragment( + cudf_fragments + FRAGMENT + ${VARIANT_NAME} + SOURCE + src/transform/jit/kernel.cu + KERNEL_INSTANCE + ${INSTANCE} + UDF_TYPE + "int(${OUTPUT_POINTERS_STR}, cudf::string_view)" + DEFINITIONS + CUDF_LTO_MODE + ARRAY_IDS + ${FRAGMENT_NAME}_FILE_INDEX + ${FRAGMENT_NAME}_INSTANCE + ARRAY_VALUES + ${FILE_INDEX} + "${INSTANCE}" + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + ) + endforeach() + endforeach() +endfunction() + +precompile_fixed_width_kernel_fragments() +precompile_string_kernel_fragments() + +rtcx_embed(cudf_fragments COMPRESSION none OUTPUT_DIRECTORY "${CMAKE_CURRENT_BINARY_DIR}/embed") # ################################################################################################## # * library targets ------------------------------------------------------------------------------- @@ -1207,7 +1314,7 @@ target_compile_options( target_include_directories( cudf PUBLIC "$" "$" - "$" + "$" PRIVATE "$" "$" "$" diff --git a/cpp/benchmarks/CMakeLists.txt b/cpp/benchmarks/CMakeLists.txt index 192ed1b978b6..4eb97a0699fd 100644 --- a/cpp/benchmarks/CMakeLists.txt +++ b/cpp/benchmarks/CMakeLists.txt @@ -399,11 +399,34 @@ ConfigureNVBench(AST_NVBENCH ast/polynomials.cpp ast/transform.cpp) # ################################################################################################## # * LTO Fragments ---------------------------------------------------------------------------- rtcx_add_embed(cudf_benchmark_fragments) -add_fragment(cudf_benchmark_fragments FRAGMENT add_f32 SOURCE binaryop/fragments/add_f32.cu) -add_fragment(cudf_benchmark_fragments FRAGMENT mul_f32 SOURCE binaryop/fragments/mul_f32.cu) +add_fragment( + cudf_benchmark_fragments + FRAGMENT + add_f32 + SOURCE + binaryop/fragments/add_f32.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} +) +add_fragment( + cudf_benchmark_fragments + FRAGMENT + mul_f32 + SOURCE + binaryop/fragments/mul_f32.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} +) rtcx_embed( - cudf_benchmark_fragments COMPRESSION none OUTPUT_DIRECTORY - "${CUDF_GENERATED_INCLUDE_DIR}/rtcx_embed" + cudf_benchmark_fragments COMPRESSION none OUTPUT_DIRECTORY "${CMAKE_CURRENT_BINARY_DIR}/embed" ) # ################################################################################################## diff --git a/cpp/cmake/Modules/AddFragment.cmake b/cpp/cmake/Modules/AddFragment.cmake index 187d9c36d594..d9cbe3538de9 100644 --- a/cpp/cmake/Modules/AddFragment.cmake +++ b/cpp/cmake/Modules/AddFragment.cmake @@ -13,8 +13,11 @@ include_guard(GLOBAL) # final library with metadata that allows it to be looked up at runtime. macro(add_fragment) set(TARGET ${ARGV0}) + set(OPTIONS) set(ONE_VALUE_ARGS FRAGMENT SOURCE KERNEL_ONLY KERNEL_INSTANCE UDF_TYPE) - set(MULTI_VALUE_ARGS DEFINITIONS ARRAY_IDS ARRAY_VALUES) + set(MULTI_VALUE_ARGS DEFINITIONS ARRAY_IDS ARRAY_VALUES INCLUDE_DIRECTORIES LINK_LIBRARIES + COMPILE_OPTIONS + ) cmake_parse_arguments(ARG "${OPTIONS}" "${ONE_VALUE_ARGS}" "${MULTI_VALUE_ARGS}" ${ARGN}) if(NOT ARG_FRAGMENT) @@ -34,7 +37,7 @@ macro(add_fragment) target_compile_options(${OBJECT_ID} PRIVATE -Xnvlink=--kernels-used=cudf_kernel_entry) endif() - set(INSTANTIATION_DIR "${CUDF_GENERATED_INCLUDE_DIR}/${TARGET}/instantiations/${ARG_FRAGMENT}") + set(INSTANTIATION_DIR "${CMAKE_CURRENT_BINARY_DIR}/${TARGET}/instantiations/${ARG_FRAGMENT}") target_include_directories(${OBJECT_ID} PRIVATE ${INSTANTIATION_DIR}) if(ARG_KERNEL_INSTANCE) @@ -53,7 +56,22 @@ macro(add_fragment) ) endif() - target_compile_definitions(${OBJECT_ID} PRIVATE CUDF_DISABLE_EXPORTS ${ARG_DEFINITIONS}) + if(ARG_INCLUDE_DIRECTORIES) + target_include_directories(${OBJECT_ID} PRIVATE ${ARG_INCLUDE_DIRECTORIES}) + endif() + + if(ARG_LINK_LIBRARIES) + target_link_libraries(${OBJECT_ID} PRIVATE ${ARG_LINK_LIBRARIES}) + endif() + + if(ARG_COMPILE_OPTIONS) + target_compile_options(${OBJECT_ID} PRIVATE ${ARG_COMPILE_OPTIONS}) + endif() + + if(ARG_DEFINITIONS) + target_compile_definitions(${OBJECT_ID} PRIVATE ${ARG_DEFINITIONS}) + endif() + set_target_properties( ${OBJECT_ID} PROPERTIES CUDA_SEPARABLE_COMPILATION ON @@ -68,17 +86,6 @@ macro(add_fragment) CUDA_STANDARD_REQUIRED ON CUDA_VISIBILITY_PRESET hidden ) - target_link_libraries( - ${OBJECT_ID} - PUBLIC CCCL::CCCL rapids_logger::rapids_logger rmm::rmm $ - PRIVATE $ $ - ZLIB::ZLIB nvcomp::nvcomp kvikio::kvikio nanoarrow::nanoarrow zstd - ) - target_include_directories( - ${OBJECT_ID} PRIVATE "$" - "$" - ) - target_compile_options(${OBJECT_ID} PRIVATE "$<$:${CUDF_CUDA_FLAGS}>") rtcx_embed_blob( ${TARGET} FILE $ DEST fragments/${ARG_FRAGMENT}.fatbin ID diff --git a/cpp/examples/string_transforms/CMakeLists.txt b/cpp/examples/string_transforms/CMakeLists.txt index c6869fe12d9e..3156e0366769 100644 --- a/cpp/examples/string_transforms/CMakeLists.txt +++ b/cpp/examples/string_transforms/CMakeLists.txt @@ -1,5 +1,5 @@ # cmake-format: off -# 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 # cmake-format: on @@ -13,7 +13,7 @@ rapids_cuda_init_architectures(string_transforms_examples) project( string_transforms_examples VERSION 0.0.1 - LANGUAGES CXX CUDA + LANGUAGES CXX CUDA ASM ) include(../fetch_dependencies.cmake) @@ -21,6 +21,16 @@ include(../fetch_dependencies.cmake) include(rapids-cmake) rapids_cmake_build_type("Release") +# Fetch librtcx and its dependencies +find_package(CUDAToolkit REQUIRED) +include(${CMAKE_CURRENT_LIST_DIR}/../../cmake/thirdparty/get_zstd.cmake) +include(${CMAKE_CURRENT_LIST_DIR}/../../cmake/thirdparty/get_xxhash.cmake) +include(${CMAKE_CURRENT_LIST_DIR}/../../cmake/thirdparty/get_nvtx.cmake) +include(${CMAKE_CURRENT_LIST_DIR}/../../cmake/thirdparty/get_rtcx.cmake) + +# Include the helper function for embedding FATBINs of fragments +include(${CMAKE_CURRENT_LIST_DIR}/../../cmake/Modules/AddFragment.cmake) + # For now, disable CMake's automatic module scanning for C++ files. There is an sccache bug in the # version RAPIDS uses in CI that causes it to handle the resulting -M* flags incorrectly with # gcc>=14. We can remove this once we upgrade to a newer sccache version. @@ -48,6 +58,32 @@ add_string_transforms_example(format_phone_precompiled format_phone_precompiled. add_string_transforms_example(localize_phone_jit localize_phone_jit.cpp) add_string_transforms_example(localize_phone_precompiled localize_phone_precompiled.cpp) -install(FILES ${CMAKE_CURRENT_LIST_DIR}/info.csv +# Compile and embed the UDFs so we can link them at runtime +rtcx_add_embed(url_log_fragments) +add_fragment( + url_log_fragments + FRAGMENT + url_component_sizes + SOURCE + url_logs/fragments.cu + DEFINITIONS + UDF_COMPUTE_SIZES + INCLUDE_DIRECTORIES + ${CMAKE_CURRENT_LIST_DIR}/url_logs + LINK_LIBRARIES + cudf::cudf +) +add_fragment( + url_log_fragments FRAGMENT url_component_output SOURCE url_logs/fragments.cu DEFINITIONS + UDF_WRITE_OUTPUT INCLUDE_DIRECTORIES ${CMAKE_CURRENT_LIST_DIR}/url_logs LINK_LIBRARIES cudf::cudf +) +rtcx_embed(url_log_fragments COMPRESSION none OUTPUT_DIRECTORY "${CMAKE_CURRENT_BINARY_DIR}/embed") + +add_string_transforms_example(url_log_transforms url_logs/transforms.cpp) +target_sources(url_log_transforms PRIVATE ${url_log_fragments_SOURCE_DIR}/url_log_fragments.s) +target_include_directories(url_log_transforms PRIVATE ${url_log_fragments_SOURCE_DIR}) +add_dependencies(url_log_transforms url_log_fragments) + +install(FILES ${CMAKE_CURRENT_LIST_DIR}/info.csv ${CMAKE_CURRENT_LIST_DIR}/url_logs/logs.csv DESTINATION bin/examples/libcudf/string_transformers ) diff --git a/cpp/examples/string_transforms/README.md b/cpp/examples/string_transforms/README.md index 1886f6d86ad8..8830e18397b6 100644 --- a/cpp/examples/string_transforms/README.md +++ b/cpp/examples/string_transforms/README.md @@ -15,6 +15,14 @@ The following examples are included: 5. `extract_email_precompiled` - Performs same transformation on the table as `output` but uses precompiled public APIs 6. `format_phone_jit` - Using a transform kernel to output a string to a pre-allocated buffer 7. `format_phone_precompiled` - Performs same transformation on the table as `preallocated` but uses precompiled public APIs +8. `url_log_transforms` - Searches raw text log lines for an embedded URL and decomposes the + first match into protocol, host, port, path, query, and fragment columns (based on https://datatracker.ietf.org/doc/html/rfc3986): + - `regex`: one six-capture `cudf::strings::extract` expression. + - `precompiled`: sequential public partition, conditional-copy, and concatenate APIs. + - `jit`: a fused byte parser compiled from CUDA source at runtime. Its sizing pass produces exact + per-row output sizes; scans create offsets and a second pass writes all six output columns. + - `lto`: the same fused parser ABI, AOT-compiled to embedded fatbins and JIT-linked with + libcudf's precompiled transform kernels. ## Compile and execute diff --git a/cpp/examples/string_transforms/url_logs/fragments.cu b/cpp/examples/string_transforms/url_logs/fragments.cu new file mode 100644 index 000000000000..7191e775f899 --- /dev/null +++ b/cpp/examples/string_transforms/url_logs/fragments.cu @@ -0,0 +1,232 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ + +#include +#include + +#include +#include + +struct range32 { + int32_t begin{}; + int32_t end{}; +}; + +struct url_ranges { + range32 protocol; + range32 host; + range32 port; + range32 path; + range32 query; + range32 fragment; +}; + +__device__ bool is_ascii_alpha(char c) { return (c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z'); } + +__device__ bool is_ascii_digit(char c) { return c >= '0' && c <= '9'; } + +__device__ bool is_hex_digit(char c) +{ + return is_ascii_digit(c) || (c >= 'A' && c <= 'F') || (c >= 'a' && c <= 'f'); +} + +// Parses the first valid URL candidate and records byte ranges for all six components. +__device__ bool parse_url(cudf::string_view input, url_ranges* out) +{ + auto n = input.size_bytes(); + *out = {}; + auto is_scheme_char = [](char c) { + return is_ascii_alpha(c) || is_ascii_digit(c) || c == '+' || c == '-' || c == '.'; + }; + auto is_unreserved = [](char c) { + return is_ascii_alpha(c) || is_ascii_digit(c) || c == '-' || c == '.' || c == '_' || c == '~'; + }; + auto is_sub_delim = [](char c) { + return c == '!' || c == '$' || c == '&' || c == '\'' || c == '(' || c == ')' || c == '*' || + c == '+' || c == ',' || c == ';' || c == '='; + }; + auto is_gen_delim = [](char c) { + return c == ':' || c == '/' || c == '?' || c == '#' || c == '[' || c == ']' || c == '@'; + }; + auto is_context_delimiter = [](char c) { + return c == ' ' || c == '\t' || c == '\n' || c == '\r' || c == '"' || c == '<' || c == '>'; + }; + + auto scheme_end = n; + for (auto i = 1; i + 2 < n; ++i) { + if (input.data()[i] == ':' && input.data()[i + 1] == '/' && input.data()[i + 2] == '/') { + scheme_end = i; + break; + } + } + if (scheme_end == n) { return false; } + + auto url_begin = scheme_end; + while (url_begin > 0 && is_scheme_char(input.data()[url_begin - 1])) { + --url_begin; + } + if (url_begin == scheme_end || !is_ascii_alpha(input.data()[url_begin])) { return false; } + + auto url_end = n; + for (auto i = scheme_end + 3; i < n; ++i) { + if (is_context_delimiter(input.data()[i])) { + url_end = i; + break; + } + } + for (auto i = url_begin; i < url_end; ++i) { + auto c = input.data()[i]; + if (c == '%') { + if (i + 2 >= url_end || !is_hex_digit(input.data()[i + 1]) || + !is_hex_digit(input.data()[i + 2])) { + return false; + } + i += 2; + } else if (!is_unreserved(c) && !is_sub_delim(c) && !is_gen_delim(c)) { + return false; + } + } + + auto hash = url_end; + for (auto i = scheme_end + 3; i < url_end; ++i) { + if (input.data()[i] == '#') { + hash = i; + break; + } + } + auto question = hash; + for (auto i = scheme_end + 3; i < hash; ++i) { + if (input.data()[i] == '?') { + question = i; + break; + } + } + auto base_end = question < hash ? question : hash; + out->protocol = {url_begin, scheme_end}; + if (question < hash) { out->query = {question + 1, hash}; } + if (hash < url_end) { out->fragment = {hash + 1, url_end}; } + + auto authority_begin = scheme_end + 3; + auto authority_end = base_end; + for (auto i = authority_begin; i < base_end; ++i) { + if (input.data()[i] == '/') { + authority_end = i; + break; + } + } + out->path = {authority_end, base_end}; + + auto host_begin = authority_begin; + for (auto i = authority_begin; i < authority_end; ++i) { + if (input.data()[i] == '@') { host_begin = i + 1; } + } + + if (host_begin < authority_end && input.data()[host_begin] == '[') { + auto close = authority_end; + for (auto i = host_begin + 1; i < authority_end; ++i) { + if (input.data()[i] == ']') { + close = i; + break; + } + } + if (close == authority_end) { return false; } + out->host = {host_begin, close + 1}; + if (close + 1 < authority_end) { + if (input.data()[close + 1] != ':') { return false; } + out->port = {close + 2, authority_end}; + } + } else { + auto colon = authority_end; + for (auto i = host_begin; i < authority_end; ++i) { + if (input.data()[i] == ':') { colon = i; } + } + out->host = {host_begin, colon}; + if (colon < authority_end) { out->port = {colon + 1, authority_end}; } + } + + for (auto i = out->port.begin; i < out->port.end; ++i) { + if (!is_ascii_digit(input.data()[i])) { return false; } + } + return true; +} + +// Computes exact output byte counts for the six URL component columns. +__device__ int compute_url_component_sizes(int32_t* protocol_size, + int32_t* host_size, + int32_t* port_size, + int32_t* path_size, + int32_t* query_size, + int32_t* fragment_size, + cudf::string_view input) +{ + *protocol_size = 0; + *host_size = 0; + *port_size = 0; + *path_size = 0; + *query_size = 0; + *fragment_size = 0; + url_ranges ranges; + if (!parse_url(input, &ranges)) { return 0; } + *protocol_size = ranges.protocol.end - ranges.protocol.begin; + *host_size = ranges.host.end - ranges.host.begin; + *port_size = ranges.port.end - ranges.port.begin; + *path_size = ranges.path.end - ranges.path.begin; + *query_size = ranges.query.end - ranges.query.begin; + *fragment_size = ranges.fragment.end - ranges.fragment.begin; + return 0; +} + +// Copies the six parsed URL components into their preallocated string buffers. +__device__ int write_url_components(cuda::std::span* protocol, + cuda::std::span* host, + cuda::std::span* port, + cuda::std::span* path, + cuda::std::span* query, + cuda::std::span* fragment, + cudf::string_view input) +{ + url_ranges ranges; + if (!parse_url(input, &ranges)) { return 0; } + cuda::std::span* outputs[] = {protocol, host, port, path, query, fragment}; + range32 components[] = { + ranges.protocol, ranges.host, ranges.port, ranges.path, ranges.query, ranges.fragment}; + for (auto component = 0; component < 6; ++component) { + auto range = components[component]; + auto size = range.end - range.begin; + if (size > 0) { memcpy(outputs[component]->data(), input.data() + range.begin, size); } + } + return 0; +} + +#ifdef UDF_COMPUTE_SIZES +// Exposes the sizing pass through the transform LTO ABI. +extern "C" __device__ int transform(int32_t* protocol_size, + int32_t* host_size, + int32_t* port_size, + int32_t* path_size, + int32_t* query_size, + int32_t* fragment_size, + cudf::string_view input) +{ + return compute_url_component_sizes( + protocol_size, host_size, port_size, path_size, query_size, fragment_size, input); +} +#else +#ifdef UDF_WRITE_OUTPUT +// Exposes the component-writing pass through the transform LTO ABI. +extern "C" __device__ int transform(cuda::std::span* protocol, + cuda::std::span* host, + cuda::std::span* port, + cuda::std::span* path, + cuda::std::span* query, + cuda::std::span* fragment, + cudf::string_view input) +{ + return write_url_components(protocol, host, port, path, query, fragment, input); +} +#else +#error "Must define either UDF_COMPUTE_SIZES or UDF_WRITE_OUTPUT" +#endif +#endif diff --git a/cpp/examples/string_transforms/url_logs/logs.csv b/cpp/examples/string_transforms/url_logs/logs.csv new file mode 100644 index 000000000000..a58d41c36bbd --- /dev/null +++ b/cpp/examples/string_transforms/url_logs/logs.csv @@ -0,0 +1,46 @@ +Timestamp,Host,Source,Level,LogLine +2026-07-31T09:14:00.000Z,edge-proxy-02,envoy,INFO,level=info service=edge method=GET status=200 url=https://api.example.com/v1/orders/123?expand=items&locale=en-GB#summary latency_ms=12 +2026-07-31T09:14:00.173Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:00.346Z,app-worker-07,checkout-worker,WARN,WARN checkout retry scheduled attempt=2 queue=payments +2026-07-31T09:14:00.519Z,identity-03,identity-api,INFO,ts=2026-07-31T09:14:00Z service=identity trace=0af765 client=198.51.100.24 target=https://login.example.com:8443/oauth2/callback?code=redacted&state=abc status=302 +2026-07-31T09:14:00.692Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:00.865Z,edge-proxy-02,envoy,INFO,edge[731]: cache=hit region=us-west-2 object=https://cdn.example.net/assets/app.8f31c2.js status=304 +2026-07-31T09:14:01.038Z,control-01,scheduler,DEBUG,level=debug component=scheduler queue_depth=0 workers=32 +2026-07-31T09:14:01.211Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:01.384Z,search-06,search-api,INFO,service=search region=ap-southeast-2 url=https://search.example.org/query?q=gpu+dataframes&page=2#results method=GET status=200 +2026-07-31T09:14:01.557Z,app-worker-07,checkout-worker,WARN,WARN checkout retry scheduled attempt=2 queue=payments +2026-07-31T09:14:01.730Z,realtime-04,realtime-gateway,INFO,INFO websocket connected client=203.0.113.71 endpoint=wss://stream.example.com/socket?token=redacted shard=4 +2026-07-31T09:14:01.903Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:02.076Z,node-17,prometheus-agent,INFO,host=node-17 process=metrics destination=http://metrics.internal.example:9090/api/v1/query?query=up result=success +2026-07-31T09:14:02.249Z,app-api-12,inventory-api,ERROR,level=error service=inventory msg=upstream_timeout attempt=3 +2026-07-31T09:14:02.422Z,control-01,scheduler,DEBUG,level=debug component=scheduler queue_depth=0 workers=32 +2026-07-31T09:14:02.595Z,graphql-02,graphql-api,INFO,request completed trace=4bf92f service=graphql method=POST url=https://api.example.com/graphql?operation=Checkout status=200 latency_ms=66 +2026-07-31T09:14:02.768Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:02.941Z,edge-proxy-02,envoy,INFO,cdn event status=206 bytes=1048576 resource=https://media.example.net/video/launch.mp4?start=120&quality=1080p region=us-west-2 +2026-07-31T09:14:03.114Z,app-worker-07,checkout-worker,WARN,WARN checkout retry scheduled attempt=2 queue=payments +2026-07-31T09:14:03.287Z,edge-proxy-01,redirector,INFO,level=info service=redirector status=301 location=http://example.com latency_ms=2 +2026-07-31T09:14:03.460Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:03.633Z,admin-01,auditd,INFO,audit actor=alice action=view target=https://admin.example.com/users/8675309#permissions result=allow +2026-07-31T09:14:03.806Z,control-01,scheduler,DEBUG,level=debug component=scheduler queue_depth=0 workers=32 +2026-07-31T09:14:03.979Z,mesh-sidecar-19,service-mesh,INFO,proxy upstream selected cluster=orders endpoint=http://orders.default.svc.cluster.local:8080/health retry=0 +2026-07-31T09:14:04.152Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:04.325Z,billing-05,billing-api,INFO,billing request tenant=42 invoice=2026-07 url=https://billing.example.com/invoices/2026-07?download=true status=200 +2026-07-31T09:14:04.498Z,app-api-12,inventory-api,ERROR,level=error service=inventory msg=upstream_timeout attempt=3 +2026-07-31T09:14:04.671Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:04.844Z,docs-01,nginx,INFO,docs access ref=homepage destination=https://docs.example.org/guides/url-parsing#query-parameters status=200 +2026-07-31T09:14:05.017Z,app-worker-07,checkout-worker,WARN,WARN checkout retry scheduled attempt=2 queue=payments +2026-07-31T09:14:05.190Z,ingest-08,event-ingest,INFO,event accepted batch=842 source=mobile callback=https://events.example.com/v3/batch?source=mobile&sdk=6.4.1 status=202 +2026-07-31T09:14:05.363Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:05.536Z,control-01,scheduler,DEBUG,level=debug component=scheduler queue_depth=0 workers=32 +2026-07-31T09:14:05.709Z,legacy-02,legacy-api,INFO,legacy request client=198.51.100.159 url=http://192.0.2.200:8000/v1/status?format=json method=GET status=200 +2026-07-31T09:14:05.882Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:06.055Z,app-worker-07,checkout-worker,WARN,WARN checkout retry scheduled attempt=2 queue=payments +2026-07-31T09:14:06.228Z,download-03,artifact-service,INFO,download complete artifact=cudf target=https://downloads.example.com/releases/package.tar.gz?signature=redacted#sha256 bytes=48234412 +2026-07-31T09:14:06.401Z,app-api-12,inventory-api,ERROR,level=error service=inventory msg=upstream_timeout attempt=3 +2026-07-31T09:14:06.574Z,k8s-worker-03,kubelet,INFO,level=info service=kube-probe msg=healthy status=200 +2026-07-31T09:14:06.747Z,control-01,scheduler,DEBUG,level=debug component=scheduler queue_depth=0 workers=32 +2026-07-31T09:14:06.920Z,edge-proxy-03,rfc-fixture,INFO,"rfc authority test url=https://user:pass@example.com:8443/a,b;c?x=/a?b#frag?part result=ok" +2026-07-31T09:14:07.093Z,directory-01,rfc-fixture,INFO,rfc ip-literal target=ldap://[2001:db8::7]/c=GB?objectClass?one result=ok +2026-07-31T09:14:07.266Z,future-net-01,rfc-fixture,INFO,rfc ipvfuture endpoint=https://[v1.fe80::a]:9443/resource result=ok +2026-07-31T09:14:07.439Z,node-17,rfc-fixture,INFO,rfc empty-authority resource=file:///etc/hosts result=ok +2026-07-31T09:14:07.612Z,edge-proxy-03,rfc-fixture,INFO,rfc pct-encoded-host url=https://example%2Ecom/a%20b result=ok diff --git a/cpp/examples/string_transforms/url_logs/transforms.cpp b/cpp/examples/string_transforms/url_logs/transforms.cpp new file mode 100644 index 000000000000..91b2420d7c42 --- /dev/null +++ b/cpp/examples/string_transforms/url_logs/transforms.cpp @@ -0,0 +1,652 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace { + +constexpr auto output_count = std::size_t{6}; + +// Shared CUDA source inserted into both runtime-compiled UDF bodies. +constexpr char parse_url_udf[] = R"***( + struct range32 { + int32_t begin{}; + int32_t end{}; + }; + struct url_ranges { + range32 protocol; + range32 host; + range32 port; + range32 path; + range32 query; + range32 fragment; + }; + // Parses the first URL candidate and records byte ranges for all six components. + auto parse_url = [&](url_ranges* out) { + *out = {}; + auto n = input.size_bytes(); + auto is_alpha = [](char c) { + return (c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z'); + }; + auto is_digit = [](char c) { return c >= '0' && c <= '9'; }; + auto is_scheme_char = [&](char c) { + return is_alpha(c) || is_digit(c) || c == '+' || c == '-' || c == '.'; + }; + auto is_hex = [&](char c) { + return is_digit(c) || (c >= 'A' && c <= 'F') || (c >= 'a' && c <= 'f'); + }; + auto is_unreserved = [&](char c) { + return is_alpha(c) || is_digit(c) || c == '-' || c == '.' || c == '_' || c == '~'; + }; + auto is_sub_delim = [](char c) { + return c == '!' || c == '$' || c == '&' || c == '\'' || c == '(' || c == ')' || c == '*' || + c == '+' || c == ',' || c == ';' || c == '='; + }; + auto is_gen_delim = [](char c) { + return c == ':' || c == '/' || c == '?' || c == '#' || c == '[' || c == ']' || c == '@'; + }; + auto is_context_delimiter = [](char c) { + return c == ' ' || c == '\t' || c == '\n' || c == '\r' || c == '"' || c == '<' || + c == '>'; + }; + + auto scheme_end = n; + for (auto i = 1; i + 2 < n; ++i) { + if (input.data()[i] == ':' && input.data()[i + 1] == '/' && input.data()[i + 2] == '/') { + scheme_end = i; + break; + } + } + if (scheme_end == n) { return false; } + + auto url_begin = scheme_end; + while (url_begin > 0 && is_scheme_char(input.data()[url_begin - 1])) { --url_begin; } + if (url_begin == scheme_end || !is_alpha(input.data()[url_begin])) { return false; } + + auto url_end = n; + for (auto i = scheme_end + 3; i < n; ++i) { + if (is_context_delimiter(input.data()[i])) { + url_end = i; + break; + } + } + for (auto i = url_begin; i < url_end; ++i) { + auto c = input.data()[i]; + if (c == '%') { + if (i + 2 >= url_end || !is_hex(input.data()[i + 1]) || !is_hex(input.data()[i + 2])) { + return false; + } + i += 2; + } else if (!is_unreserved(c) && !is_sub_delim(c) && !is_gen_delim(c)) { + return false; + } + } + + auto hash = url_end; + for (auto i = scheme_end + 3; i < url_end; ++i) { + if (input.data()[i] == '#') { + hash = i; + break; + } + } + auto question = hash; + for (auto i = scheme_end + 3; i < hash; ++i) { + if (input.data()[i] == '?') { + question = i; + break; + } + } + auto base_end = question < hash ? question : hash; + out->protocol = {url_begin, scheme_end}; + if (question < hash) { out->query = {question + 1, hash}; } + if (hash < url_end) { out->fragment = {hash + 1, url_end}; } + + auto authority_begin = scheme_end + 3; + auto authority_end = base_end; + for (auto i = authority_begin; i < base_end; ++i) { + if (input.data()[i] == '/') { + authority_end = i; + break; + } + } + out->path = {authority_end, base_end}; + + auto host_begin = authority_begin; + for (auto i = authority_begin; i < authority_end; ++i) { + if (input.data()[i] == '@') { host_begin = i + 1; } + } + if (host_begin < authority_end && input.data()[host_begin] == '[') { + auto close = authority_end; + for (auto i = host_begin + 1; i < authority_end; ++i) { + if (input.data()[i] == ']') { + close = i; + break; + } + } + if (close == authority_end) { return false; } + out->host = {host_begin, close + 1}; + if (close + 1 < authority_end) { + if (input.data()[close + 1] != ':') { return false; } + out->port = {close + 2, authority_end}; + } + } else { + auto colon = authority_end; + for (auto i = host_begin; i < authority_end; ++i) { + if (input.data()[i] == ':') { colon = i; } + } + out->host = {host_begin, colon}; + if (colon < authority_end) { out->port = {colon + 1, authority_end}; } + } + for (auto i = out->port.begin; i < out->port.end; ++i) { + if (!is_digit(input.data()[i])) { return false; } + } + return true; + }; +)***"; + +// Builds the sizing UDF by inserting the shared parser into a self-contained device function. +std::string const url_component_sizes_udf = std::string{R"***( +// Computes exact output byte counts for the six URL component columns. +__device__ int compute_url_component_sizes(int32_t* protocol_size, + int32_t* host_size, + int32_t* port_size, + int32_t* path_size, + int32_t* query_size, + int32_t* fragment_size, + cudf::string_view input) { + *protocol_size = *host_size = *port_size = 0; + *path_size = *query_size = *fragment_size = 0; +)***"} + parse_url_udf + R"***( + url_ranges ranges; + if (!parse_url(&ranges)) { return 0; } + *protocol_size = ranges.protocol.end - ranges.protocol.begin; + *host_size = ranges.host.end - ranges.host.begin; + *port_size = ranges.port.end - ranges.port.begin; + *path_size = ranges.path.end - ranges.path.begin; + *query_size = ranges.query.end - ranges.query.begin; + *fragment_size = ranges.fragment.end - ranges.fragment.begin; + return 0; +} +)***"; + +// Builds the output UDF from the same parser so both CUDA passes use identical ranges. +std::string const url_component_output_udf = std::string{R"***( +// Copies the six parsed URL components into their preallocated string buffers. +__device__ int write_url_components(cuda::std::span* protocol, + cuda::std::span* host, + cuda::std::span* port, + cuda::std::span* path, + cuda::std::span* query, + cuda::std::span* fragment, + cudf::string_view input) { +)***"} + parse_url_udf + R"***( + url_ranges ranges; + if (!parse_url(&ranges)) { return 0; } + cuda::std::span* outputs[] = {protocol, host, port, path, query, fragment}; + range32 components[] = { + ranges.protocol, ranges.host, ranges.port, ranges.path, ranges.query, ranges.fragment}; + for (auto component = 0; component < 6; ++component) { + auto range = components[component]; + auto size = range.end - range.begin; + if (size > 0) { memcpy(outputs[component]->data(), input.data() + range.begin, size); } + } + return 0; +} +)***"; + +constexpr std::string_view usage = + "usage: url_log_transforms INPUT.csv OUTPUT.csv ROWS\n" + " url_log_transforms INPUT.csv OUTPUT.csv ROWS " + "<--warm|--cold|--cold-warm-pch>\n" + " url_log_transforms \n"; + +// warmup the PCH cache +void warmup_pch(cudf::column_view input, + rmm::cuda_stream_view stream, + rmm::device_async_resource_ref mr) +{ + constexpr char udf[] = R"***( +__device__ int transform(int32_t* output, cudf::string_view input) { + *output = input.size_bytes(); + return 0; +} +)***"; + cudf::transform_input inputs[] = {input}; + cudf::transform_output const output{cudf::data_type{cudf::type_id::INT32}, + cudf::output_nullability::ALL_VALID}; + std::vector const outputs{output}; + auto result = cudf::transform(udf, + cudf::udf_source_type::CUDA, + cudf::null_aware::NO, + std::nullopt, + inputs, + outputs, + {}, + std::nullopt, + stream, + mr); + stream.synchronize(); +} + +// Extracts RFC 3986-style hierarchical URI components from unstructured log lines. +[[nodiscard]] std::unique_ptr run_regex(cudf::column_view input, + rmm::cuda_stream_view stream, + rmm::device_async_resource_ref mr) +{ + // Derived from RFC 3986 Appendix B (https://www.rfc-editor.org/info/rfc3986/#page-50). The + // authority capture is expanded into optional userinfo plus host and port, and Appendix C + // delimiters bound the URI within a log line. + static auto program = cudf::strings::regex_program::create( + R"((?:^|[^A-Za-z0-9+.-])([A-Za-z][A-Za-z0-9+.-]*):\/\/(?:[^@\/?# \t\n\r"<>]*@)?(\[[^\]\/?# \t\n\r"<>]*\]|[^\/:?# \t\n\r"<>]*)(?::([0-9]*))?([^?# \t\n\r"<>]*)(?:\?([^# \t\n\r"<>]*))?(?:#([^ \t\n\r"<>]*))?(?:$|[ \t\n\r"<>]))"); + auto extracted = cudf::strings::extract(cudf::strings_column_view{input}, *program, stream, mr); + auto columns = extracted->release(); + auto empty = cudf::string_scalar{"", true, stream, mr}; + for (auto& column : columns) { + column = cudf::replace_nulls(column->view(), empty, stream, mr); + } + return std::make_unique(std::move(columns)); +} + +// Decomposes key-value URL tokens using only precompiled libcudf string primitives. +[[nodiscard]] std::unique_ptr run_precompiled(cudf::column_view input, + rmm::cuda_stream_view stream, + rmm::device_async_resource_ref mr) +{ + // Materialize the delimiters used by each partitioning stage. + auto empty = cudf::string_scalar{"", true, stream, mr}; + auto scheme_separator = cudf::string_scalar{"://", true, stream, mr}; + auto marker_separator = cudf::string_scalar{"=", true, stream, mr}; + auto token_separator = cudf::string_scalar{" ", true, stream, mr}; + auto hash = cudf::string_scalar{"#", true, stream, mr}; + auto question = cudf::string_scalar{"?", true, stream, mr}; + auto slash = cudf::string_scalar{"/", true, stream, mr}; + auto at = cudf::string_scalar{"@", true, stream, mr}; + auto right_bracket = cudf::string_scalar{"]", true, stream, mr}; + auto left_bracket = cudf::string_scalar{"[", true, stream, mr}; + auto colon = cudf::string_scalar{":", true, stream, mr}; + + // Mark rows containing an authority-style URI and split at the first "://". + auto has_url = + cudf::strings::contains(cudf::strings_column_view{input}, scheme_separator, stream, mr); + auto scheme_table = + cudf::strings::partition(cudf::strings_column_view{input}, scheme_separator, stream, mr); + auto scheme_columns = scheme_table->release(); + + // Extract the scheme from the key-value token immediately preceding "://". + auto marker_table = cudf::strings::rpartition( + cudf::strings_column_view{scheme_columns[0]->view()}, marker_separator, stream, mr); + auto marker_columns = marker_table->release(); + + // Stop at the first space so later log fields are excluded from the URI. + auto token_table = cudf::strings::partition( + cudf::strings_column_view{scheme_columns[2]->view()}, token_separator, stream, mr); + auto token_columns = token_table->release(); + + // Split off the fragment; everything after the first '#' belongs to it. + auto fragment_table = + cudf::strings::partition(cudf::strings_column_view{token_columns[0]->view()}, hash, stream, mr); + auto fragment_columns = fragment_table->release(); + + // Split the pre-fragment portion at the first '?' to isolate the query. + auto query_table = cudf::strings::partition( + cudf::strings_column_view{fragment_columns[0]->view()}, question, stream, mr); + auto query_columns = query_table->release(); + + // Split the remaining hierarchical part at its first slash into authority and path. + auto authority_path_table = cudf::strings::partition( + cudf::strings_column_view{query_columns[0]->view()}, slash, stream, mr); + auto authority_path_columns = authority_path_table->release(); + + // Reattach the slash delimiter to produce the RFC path value. + auto path = cudf::strings::concatenate( + cudf::table_view{{authority_path_columns[1]->view(), authority_path_columns[2]->view()}}, + empty, + cudf::string_scalar{"", false, stream, mr}, + cudf::strings::separator_on_nulls::YES, + stream, + mr); + + // Remove optional userinfo by retaining everything after the authority's last '@'. + auto has_userinfo = cudf::strings::contains( + cudf::strings_column_view{authority_path_columns[0]->view()}, at, stream, mr); + auto userinfo_table = cudf::strings::rpartition( + cudf::strings_column_view{authority_path_columns[0]->view()}, at, stream, mr); + auto userinfo_columns = userinfo_table->release(); + auto host_port = cudf::copy_if_else(userinfo_columns[2]->view(), + authority_path_columns[0]->view(), + has_userinfo->view(), + stream, + mr); + + // Bracketed IP literals and regular hosts require different port splitting rules. + auto is_ip_literal = cudf::strings::starts_with( + cudf::strings_column_view{host_port->view()}, left_bracket, stream, mr); + auto bracket_table = cudf::strings::partition( + cudf::strings_column_view{host_port->view()}, right_bracket, stream, mr); + auto bracket_columns = bracket_table->release(); + + // Preserve both brackets as part of an IP-literal host. + auto bracket_host = cudf::strings::concatenate( + cudf::table_view{{bracket_columns[0]->view(), bracket_columns[1]->view()}}, + empty, + cudf::string_scalar{"", false, stream, mr}, + cudf::strings::separator_on_nulls::YES, + stream, + mr); + + // For an IP literal, parse an optional port only after the closing bracket. + auto bracket_port_table = cudf::strings::partition( + cudf::strings_column_view{bracket_columns[2]->view()}, colon, stream, mr); + auto bracket_port_columns = bracket_port_table->release(); + + // For a regular authority, treat the final colon as the port separator. + auto has_regular_port = + cudf::strings::contains(cudf::strings_column_view{host_port->view()}, colon, stream, mr); + auto regular_table = + cudf::strings::rpartition(cudf::strings_column_view{host_port->view()}, colon, stream, mr); + auto regular_columns = regular_table->release(); + auto regular_host = cudf::copy_if_else( + regular_columns[0]->view(), host_port->view(), has_regular_port->view(), stream, mr); + auto regular_port = + cudf::copy_if_else(regular_columns[2]->view(), empty, has_regular_port->view(), stream, mr); + + // Select the bracketed or regular host/port result for each row. + auto host = cudf::copy_if_else( + bracket_host->view(), regular_host->view(), is_ip_literal->view(), stream, mr); + auto port = cudf::copy_if_else( + bracket_port_columns[2]->view(), regular_port->view(), is_ip_literal->view(), stream, mr); + + // Convert missing components to empty strings and blank rows without a URL. + auto normalize = [&](cudf::column_view column) { + auto no_nulls = cudf::replace_nulls(column, empty, stream, mr); + return cudf::copy_if_else(no_nulls->view(), empty, has_url->view(), stream, mr); + }; + + // Return the six columns in the same order used by the regex and CUDA implementations. + std::vector> result; + result.reserve(output_count); + result.push_back(normalize(marker_columns[2]->view())); + result.push_back(normalize(host->view())); + result.push_back(normalize(port->view())); + result.push_back(normalize(path->view())); + result.push_back(normalize(query_columns[2]->view())); + result.push_back(normalize(fragment_columns[2]->view())); + return std::make_unique(std::move(result)); +} + +// Runs either the runtime-compiled CUDA-string UDFs or their AOT fatbin/LTO counterparts. +[[nodiscard]] std::unique_ptr run_jit(cudf::column_view input, + bool use_lto, + rmm::cuda_stream_view stream, + rmm::device_async_resource_ref mr) +{ + cudf::transform_output const size_spec{cudf::data_type{cudf::type_id::INT32}, + cudf::output_nullability::ALL_VALID}; + std::vector const size_outputs(output_count, size_spec); + cudf::transform_input inputs[] = {input}; + std::unique_ptr sizes; + + if (use_lto) { + auto range = url_log_fragments::file_ranges[url_log_fragments::url_component_sizes]; + auto fragment = url_log_fragments::files.subspan(range[0], range[1]); + sizes = cudf::transform_lto(fragment, + cudf::lto_binary_type::FATBIN, + cudf::null_aware::NO, + std::nullopt, + inputs, + size_outputs, + {}, + std::nullopt, + stream, + mr); + } else { + sizes = cudf::transform(url_component_sizes_udf, + cudf::udf_source_type::CUDA, + cudf::null_aware::NO, + std::nullopt, + inputs, + size_outputs, + {}, + std::nullopt, + stream, + mr); + } + + std::vector> offsets; + offsets.reserve(output_count); + for (auto& string_sizes : sizes->view()) { + auto run_ends = cudf::scan(string_sizes, + *cudf::make_sum_aggregation(), + cudf::scan_type::INCLUSIVE, + cudf::null_policy::EXCLUDE, + stream, + mr); + auto zero = cudf::numeric_scalar{0, true, stream, mr}; + auto first = cudf::make_column_from_scalar(zero, 1, stream, mr); + offsets.push_back(cudf::concatenate( + std::vector{first->view(), run_ends->view()}, stream, mr)); + } + + cudf::transform_output const output_spec{cudf::data_type{cudf::type_id::STRING}, + cudf::output_nullability::ALL_VALID}; + std::vector const outputs(output_count, output_spec); + if (use_lto) { + auto range = url_log_fragments::file_ranges[url_log_fragments::url_component_output]; + auto fragment = url_log_fragments::files.subspan(range[0], range[1]); + return cudf::transform_lto(fragment, + cudf::lto_binary_type::FATBIN, + cudf::null_aware::NO, + std::nullopt, + inputs, + outputs, + std::move(offsets), + std::nullopt, + stream, + mr); + } + return cudf::transform(url_component_output_udf, + cudf::udf_source_type::CUDA, + cudf::null_aware::NO, + std::nullopt, + inputs, + outputs, + std::move(offsets), + std::nullopt, + stream, + mr); +} + +} // namespace + +int main(int argc, char const** argv) +try { + if (argc == 2 && + (std::string_view{argv[1]} == "--help" || std::string_view{argv[1]} == "usage")) { + std::cout << usage; + return EXIT_SUCCESS; + } + if (argc != 5 && argc != 6) { + throw std::invalid_argument("invalid arguments; run url_log_transforms --help for usage"); + } + + auto input_path = std::string{argv[1]}; + auto output_path = std::string{argv[2]}; + auto impl = std::string_view{argv[3]}; + if (impl != "regex" && impl != "precompiled" && impl != "cuda-jit" && impl != "lto-jit") { + throw std::invalid_argument("executor must be regex, precompiled, cuda-jit, or lto-jit"); + } + auto requested_rows = std::stoll(argv[4]); + auto is_jit = impl == "cuda-jit" || impl == "lto-jit"; + if (is_jit && argc != 6) { + throw std::invalid_argument("cuda-jit and lto-jit require a warm-up control"); + } + if (!is_jit && argc != 5) { + throw std::invalid_argument("regex and precompiled do not accept a warm-up control"); + } + auto warmup_control = argc == 6 ? std::string_view{argv[5]} : std::string_view{"none"}; + if (is_jit && warmup_control != "--warm" && warmup_control != "--cold" && + warmup_control != "--cold-warm-pch") { + throw std::invalid_argument("warm-up control must be --warm, --cold, or --cold-warm-pch"); + } + if (requested_rows < 0 || requested_rows > std::numeric_limits::max()) { + throw std::invalid_argument("ROWS is outside the cudf::size_type range"); + } + nvtxRangePush("url_log_process"); + auto process_start = std::chrono::steady_clock::now(); + auto rows = static_cast(requested_rows); + auto use_lto = impl == "lto-jit"; + auto stream = cudf::get_default_stream(); + auto upstream_mr = cudf::get_current_device_resource_ref(); + // Tracks setup, measured work, and output + rmm::mr::statistics_resource_adaptor whole_stats{upstream_mr}; + auto whole_mr = rmm::device_async_resource_ref{whole_stats}; + cudf::set_current_device_resource(whole_mr); + + nvtxRangePush("url_log_setup"); + auto read_options = cudf::io::csv_reader_options::builder(cudf::io::source_info{input_path}) + .header(0) + .use_cols_names({"LogLine"}) + .build(); + auto input = cudf::io::read_csv(read_options).tbl; + if (rows != input->num_rows()) { + input = + cudf::sample(input->view(), rows, cudf::sample_with_replacement::TRUE, 0, stream, whole_mr); + } + stream.synchronize(); + auto input_view = input->get_column(0).view(); + auto logical_input_bytes = cudf::strings_column_view{input_view}.chars_size(stream); + nvtxRangePop(); + + // Tracks measured work; nested allocations also update whole_stats. + rmm::mr::statistics_resource_adaptor measured_stats{whole_mr}; + auto measured_mr = rmm::device_async_resource_ref{measured_stats}; + auto run_transform = [&](rmm::device_async_resource_ref mr) { + if (impl == "regex") { + return run_regex(input_view, stream, mr); + } else if (impl == "precompiled") { + return run_precompiled(input_view, stream, mr); + } else { + return run_jit(input_view, use_lto, stream, mr); + } + }; + + std::unique_ptr result; + auto warmup_duration = std::chrono::steady_clock::duration::zero(); + + if (warmup_control == "--cold-warm-pch") { + // Do not track warm-up allocations. + cudf::set_current_device_resource(upstream_mr); + nvtxRangePush("url_log_warmup"); + warmup_pch(input_view, stream, upstream_mr); + nvtxRangePop(); + cudf::set_current_device_resource(whole_mr); + } else if (warmup_control == "--warm") { + // Do not track warm-up allocations. + cudf::set_current_device_resource(upstream_mr); + stream.synchronize(); + auto warmup_start = std::chrono::steady_clock::now(); + nvtxRangePush("url_log_warmup"); + result = run_transform(upstream_mr); + stream.synchronize(); + nvtxRangePop(); + warmup_duration = std::chrono::steady_clock::now() - warmup_start; + result.reset(); + cudf::set_current_device_resource(whole_mr); + } + + // Measured allocations update both statistics scopes. + cudf::set_current_device_resource(measured_mr); + stream.synchronize(); + auto measured_start = std::chrono::steady_clock::now(); + nvtxRangePush("url_log_measured"); + result = run_transform(measured_mr); + stream.synchronize(); + nvtxRangePop(); + auto measured_duration = std::chrono::steady_clock::now() - measured_start; + + if (output_path != "-") { + // Exclude output serialization from measured statistics. + cudf::set_current_device_resource(whole_mr); + auto write_options = + cudf::io::csv_writer_options::builder(cudf::io::sink_info{output_path}, result->view()) + .include_header(true) + .names({"protocol", "host", "port", "path", "query", "fragment"}) + .build(); + cudf::io::write_csv(write_options); + } + + // Read measured and broader workload scopes separately. + auto measured_bytes = measured_stats.get_bytes_counter(); + auto whole_bytes = whole_stats.get_bytes_counter(); + auto output_allocated_bytes = result->alloc_size(); + auto input_gib = static_cast(logical_input_bytes) / static_cast(1ULL << 30); + auto whole_duration = std::chrono::steady_clock::now() - process_start; + auto warmup_seconds = std::chrono::duration{warmup_duration}.count(); + auto measured_seconds = std::chrono::duration{measured_duration}.count(); + auto whole_seconds = std::chrono::duration{whole_duration}.count(); + std::cout << std::format( + "executor={}\nwarmup_control={}\nrows={}\nwarmup_seconds={}\n" + "measured_cpu_wall_seconds={}\nrows_per_second={}\n" + "input_gib_per_second={}\nlogical_input_bytes={}\noutput_allocated_bytes={}\n" + "peak_memory_bytes={}\n" + "total_allocated_bytes={}\nallocated_bytes_per_call={}\nmeasured_gpu_peak_bytes={}\n" + "measured_gpu_allocation_volume_bytes={}\nwhole_workload_seconds={}\n" + "whole_gpu_peak_bytes={}\nwhole_gpu_allocation_volume_bytes={}\n", + impl, + warmup_control, + rows, + warmup_seconds, + measured_seconds, + static_cast(rows) / measured_seconds, + input_gib / measured_seconds, + logical_input_bytes, + output_allocated_bytes, + measured_bytes.peak, + measured_bytes.total, + measured_bytes.total, + measured_bytes.peak, + measured_bytes.total, + whole_seconds, + whole_bytes.peak, + whole_bytes.total); + result.reset(); + input.reset(); + cudf::set_current_device_resource(upstream_mr); + nvtxRangePop(); + return EXIT_SUCCESS; +} catch (std::exception const& error) { + std::cerr << error.what() << '\n'; + return EXIT_FAILURE; +} diff --git a/cpp/src/jit/cache.hpp b/cpp/src/jit/cache.hpp index f96d58c587af..cd0a68a6c7ef 100644 --- a/cpp/src/jit/cache.hpp +++ b/cpp/src/jit/cache.hpp @@ -64,7 +64,7 @@ struct [[nodiscard]] kernel { rtcx::cuda_dim3 block_dim, uint32_t shared_mem_bytes, cuda::stream_ref stream, - Args&&... args) + Args&&... args) const requires(sizeof...(Args) > 0) { void const* params[] = {&args...}; // NOLINT(modernize-avoid-c-arrays) diff --git a/cpp/src/transform/transform.cu b/cpp/src/transform/transform.cu index a1fc1bca873d..8c3662c58816 100644 --- a/cpp/src/transform/transform.cu +++ b/cpp/src/transform/transform.cu @@ -389,8 +389,7 @@ std::string reflect_udf_signature(bool is_null_aware, { std::vector in_types; - for (size_t i = 0; i < inputs.size(); i++) { - auto& in = inputs[i]; + for (auto& in : inputs) { auto element = std::visit([&](auto& c) { return reflect_input_element(c, use_physical_types); }, in); in_types.push_back(is_null_aware ? std::format("cuda::std::optional<{}>", element) : element); @@ -398,8 +397,7 @@ std::string reflect_udf_signature(bool is_null_aware, std::vector out_types; - for (size_t i = 0; i < outputs.size(); i++) { - auto& out = outputs[i]; + for (auto& out : outputs) { auto element = std::visit([&](auto& c) { return reflect_output_element(c, use_physical_types); }, out); out_types.push_back(is_null_aware ? std::format("cuda::std::optional<{}> *", element) @@ -407,16 +405,20 @@ std::string reflect_udf_signature(bool is_null_aware, } std::vector params; - if (has_user_data) { params.push_back("void*"); } + if (has_user_data) { + params.emplace_back("void*"); + params.emplace_back("cudf::size_type"); + } params.insert(params.end(), out_types.begin(), out_types.end()); params.insert(params.end(), in_types.begin(), in_types.end()); auto joined = params.empty() ? "" - : std::accumulate(std::next(params.begin()), params.end(), params[0], [](auto a, auto b) { - return std::format("{}, {}", a, b); - }); + : std::accumulate( + std::next(params.begin()), params.end(), params[0], [](auto const& a, auto const& b) { + return std::format("{}, {}", a, b); + }); return std::format("int({})", joined); } diff --git a/cpp/tests/CMakeLists.txt b/cpp/tests/CMakeLists.txt index d268ed76d3b2..5d0387a04b25 100644 --- a/cpp/tests/CMakeLists.txt +++ b/cpp/tests/CMakeLists.txt @@ -695,23 +695,97 @@ ConfigureTest(AST_TEST ast/transform_tests.cpp ast/ast_tree_tests.cpp ast/jit_ex rtcx_add_embed(cudf_test_fragments) add_fragment( - cudf_test_fragments FRAGMENT bankers_rounding SOURCE transform/fragments/bankers_rounding.cu + cudf_test_fragments + LINK_CUDF_DEPS + FRAGMENT + bankers_rounding + SOURCE + transform/fragments/bankers_rounding.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} ) -add_fragment(cudf_test_fragments FRAGMENT distance SOURCE transform/fragments/distance.cu) +add_fragment( + cudf_test_fragments + LINK_CUDF_DEPS + FRAGMENT + distance + SOURCE + transform/fragments/distance.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} +) -add_fragment(cudf_test_fragments FRAGMENT invsqrt SOURCE transform/fragments/invsqrt.cu) +add_fragment( + cudf_test_fragments + LINK_CUDF_DEPS + FRAGMENT + invsqrt + SOURCE + transform/fragments/invsqrt.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} +) -add_fragment(cudf_test_fragments FRAGMENT lehmer_mean SOURCE transform/fragments/lehmer_mean.cu) +add_fragment( + cudf_test_fragments + LINK_CUDF_DEPS + FRAGMENT + lehmer_mean + SOURCE + transform/fragments/lehmer_mean.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} +) add_fragment( - cudf_test_fragments FRAGMENT sum_of_squares SOURCE transform/fragments/sum_of_squares.cu + cudf_test_fragments + LINK_CUDF_DEPS + FRAGMENT + sum_of_squares + SOURCE + transform/fragments/sum_of_squares.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} ) -add_fragment(cudf_test_fragments FRAGMENT to_upper SOURCE transform/fragments/to_upper.cu) +add_fragment( + cudf_test_fragments + LINK_CUDF_DEPS + FRAGMENT + to_upper + SOURCE + transform/fragments/to_upper.cu + LINK_LIBRARIES + ${LIBCUDF_FRAGMENT_LINK_LIBRARIES} + INCLUDE_DIRECTORIES + ${LIBCUDF_FRAGMENT_INCLUDE_DIRECTORIES} + COMPILE_OPTIONS + ${LIBCUDF_FRAGMENT_COMPILE_OPTIONS} +) rtcx_embed( - cudf_test_fragments COMPRESSION none OUTPUT_DIRECTORY "${CUDF_GENERATED_INCLUDE_DIR}/rtcx_embed" + cudf_test_fragments COMPRESSION none OUTPUT_DIRECTORY "${CMAKE_CURRENT_BINARY_DIR}/embed" ) ConfigureTest(