From 1979d826657a739b0fe477f570abf70d0cca0844 Mon Sep 17 00:00:00 2001 From: haixuantao Date: Mon, 13 Apr 2026 21:12:43 +0200 Subject: [PATCH] feat: Add cmake-arrowcuda and rust-cuda dataflow examples Add two CUDA IPC examples demonstrating Arrow CUDA data transfer between nodes: - cmake-arrowcuda-dataflow: C++ node receiving Arrow CUDA IPC buffers - rust-cuda-dataflow: Rust node receiving Arrow CUDA IPC buffers Co-Authored-By: Claude Opus 4.6 (1M context) --- Cargo.toml | 1 + examples/cmake-arrowcuda-dataflow/.gitignore | 2 + .../cmake-arrowcuda-dataflow/CMakeLists.txt | 33 ++++ .../DoraTargets.cmake | 107 +++++++++++ .../cmake-arrowcuda-dataflow/dataflow.yml | 14 ++ .../cmake-arrowcuda-dataflow/node/main.cc | 178 ++++++++++++++++++ examples/cmake-arrowcuda-dataflow/sender.py | 19 ++ examples/rust-cuda-dataflow/dataflow.yml | 15 ++ .../rust-cuda-dataflow/receiver/Cargo.toml | 10 + .../rust-cuda-dataflow/receiver/src/main.rs | 95 ++++++++++ examples/rust-cuda-dataflow/sender.py | 19 ++ 11 files changed, 493 insertions(+) create mode 100644 examples/cmake-arrowcuda-dataflow/.gitignore create mode 100644 examples/cmake-arrowcuda-dataflow/CMakeLists.txt create mode 100644 examples/cmake-arrowcuda-dataflow/DoraTargets.cmake create mode 100644 examples/cmake-arrowcuda-dataflow/dataflow.yml create mode 100644 examples/cmake-arrowcuda-dataflow/node/main.cc create mode 100644 examples/cmake-arrowcuda-dataflow/sender.py create mode 100644 examples/rust-cuda-dataflow/dataflow.yml create mode 100644 examples/rust-cuda-dataflow/receiver/Cargo.toml create mode 100644 examples/rust-cuda-dataflow/receiver/src/main.rs create mode 100644 examples/rust-cuda-dataflow/sender.py diff --git a/Cargo.toml b/Cargo.toml index 6fcdb27..0f5366f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,6 +5,7 @@ members = [ "nodes/sink-node", "nodes/status-node", "nodes/sink-dynamic-node", + "examples/rust-cuda-dataflow/receiver", ] [package] diff --git a/examples/cmake-arrowcuda-dataflow/.gitignore b/examples/cmake-arrowcuda-dataflow/.gitignore new file mode 100644 index 0000000..3c0160d --- /dev/null +++ b/examples/cmake-arrowcuda-dataflow/.gitignore @@ -0,0 +1,2 @@ +build/ +out/ diff --git a/examples/cmake-arrowcuda-dataflow/CMakeLists.txt b/examples/cmake-arrowcuda-dataflow/CMakeLists.txt new file mode 100644 index 0000000..15f0779 --- /dev/null +++ b/examples/cmake-arrowcuda-dataflow/CMakeLists.txt @@ -0,0 +1,33 @@ +cmake_minimum_required(VERSION 3.15) +project(cmake-arrowcuda-dataflow LANGUAGES CXX) + +set(CMAKE_CXX_STANDARD 20) +set(CMAKE_CXX_FLAGS "-fPIC") + +find_package(Arrow REQUIRED) +find_package(ArrowCUDA REQUIRED) + +include(DoraTargets.cmake) + +link_directories(${dora_link_dirs}) + +add_executable(node node/main.cc ${node_bridge}) +add_dependencies(node Dora_cxx) +target_include_directories(node PRIVATE ${dora_cxx_include_dir}) +target_link_libraries(node dora_node_api_cxx Arrow::arrow_shared ArrowCUDA::arrow_cuda_shared) + +# Platform-specific system libraries required by the dora static library +if(APPLE) + target_link_libraries(node + "-framework CoreServices" "-framework CoreFoundation" "-framework Security" + pthread m c resolv) +elseif(UNIX) + target_link_libraries(node pthread dl m rt) +elseif(WIN32) + target_link_libraries(node + advapi32 userenv kernel32 ws2_32 bcrypt ncrypt schannel ntdll iphlpapi + cfgmgr32 credui crypt32 cryptnet fwpuclnt gdi32 msimg32 mswsock ole32 + oleaut32 opengl32 secur32 shell32 synchronization user32 winspool) +endif() + +install(TARGETS node DESTINATION ${CMAKE_CURRENT_SOURCE_DIR}/build) diff --git a/examples/cmake-arrowcuda-dataflow/DoraTargets.cmake b/examples/cmake-arrowcuda-dataflow/DoraTargets.cmake new file mode 100644 index 0000000..16cdc4a --- /dev/null +++ b/examples/cmake-arrowcuda-dataflow/DoraTargets.cmake @@ -0,0 +1,107 @@ +set(DORA_ROOT_DIR "" CACHE FILEPATH "Path to the root of dora") + +set(dora_c_include_dir "${CMAKE_CURRENT_BINARY_DIR}/include/c") + +set(dora_cxx_include_dir "${CMAKE_CURRENT_BINARY_DIR}/include/cxx") +set(node_bridge "${CMAKE_CURRENT_BINARY_DIR}/node_bridge.cc") +set(operator_bridge "${CMAKE_CURRENT_BINARY_DIR}/operator_bridge.cc") + +if(DORA_ROOT_DIR) + include(ExternalProject) + ExternalProject_Add(Dora + SOURCE_DIR ${DORA_ROOT_DIR} + BUILD_IN_SOURCE True + CONFIGURE_COMMAND "" + BUILD_COMMAND + cargo build + --package dora-node-api-c + && + cargo build + --package dora-operator-api-c + && + cargo build + --package dora-node-api-cxx + && + cargo build + --package dora-operator-api-cxx + INSTALL_COMMAND "" + ) + + add_custom_command(OUTPUT ${node_bridge} ${dora_cxx_include_dir} ${operator_bridge} ${dora_c_include_dir} + WORKING_DIRECTORY ${DORA_ROOT_DIR} + DEPENDS Dora + COMMAND + mkdir -p ${dora_cxx_include_dir} + && + mkdir -p ${CMAKE_CURRENT_BINARY_DIR}/include/c + && + cp target/cxxbridge/dora-node-api-cxx/src/lib.rs.cc ${node_bridge} + && + cp target/cxxbridge/dora-node-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-node-api.h + && + cp target/cxxbridge/dora-operator-api-cxx/src/lib.rs.cc ${operator_bridge} + && + cp target/cxxbridge/dora-operator-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-operator-api.h + && + cp -r apis/c/node ${CMAKE_CURRENT_BINARY_DIR}/include/c + && + cp -r apis/c/operator ${CMAKE_CURRENT_BINARY_DIR}/include/c + + ) + + add_custom_target(Dora_c DEPENDS ${dora_c_include_dir}) + add_custom_target(Dora_cxx DEPENDS ${node_bridge} ${operator_bridge} ${dora_cxx_include_dir}) + set(dora_link_dirs ${DORA_ROOT_DIR}/target/debug) +else() + include(ExternalProject) + ExternalProject_Add(Dora + PREFIX ${CMAKE_CURRENT_BINARY_DIR}/dora + GIT_REPOSITORY https://github.com/dora-rs/dora.git + GIT_TAG main + BUILD_IN_SOURCE True + CONFIGURE_COMMAND "" + BUILD_COMMAND + cargo build + --package dora-node-api-c + --target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target + && + cargo build + --package dora-operator-api-c + --target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target + && + cargo build + --package dora-node-api-cxx + --target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target + && + cargo build + --package dora-operator-api-cxx + --target-dir ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target + INSTALL_COMMAND "" + ) + + add_custom_command(OUTPUT ${node_bridge} ${dora_cxx_include_dir} ${operator_bridge} ${dora_c_include_dir} + WORKING_DIRECTORY ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target + DEPENDS Dora + COMMAND + mkdir -p ${CMAKE_CURRENT_BINARY_DIR}/include/c + && + mkdir -p ${dora_cxx_include_dir} + && + cp cxxbridge/dora-node-api-cxx/src/lib.rs.cc ${node_bridge} + && + cp cxxbridge/dora-node-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-node-api.h + && + cp cxxbridge/dora-operator-api-cxx/src/lib.rs.cc ${operator_bridge} + && + cp cxxbridge/dora-operator-api-cxx/src/lib.rs.h ${dora_cxx_include_dir}/dora-operator-api.h + && + cp ../apis/c/node ${CMAKE_CURRENT_BINARY_DIR}/include/c -r + && + cp ../apis/c/operator ${CMAKE_CURRENT_BINARY_DIR}/include/c -r + ) + + set(dora_link_dirs ${CMAKE_CURRENT_BINARY_DIR}/dora/src/Dora/target/debug) + + add_custom_target(Dora_c DEPENDS ${dora_c_include_dir}) + add_custom_target(Dora_cxx DEPENDS ${node_bridge} ${operator_bridge} ${dora_cxx_include_dir}) +endif() diff --git a/examples/cmake-arrowcuda-dataflow/dataflow.yml b/examples/cmake-arrowcuda-dataflow/dataflow.yml new file mode 100644 index 0000000..31869d5 --- /dev/null +++ b/examples/cmake-arrowcuda-dataflow/dataflow.yml @@ -0,0 +1,14 @@ +nodes: + - id: sender + path: sender.py + inputs: + next: receiver/next + outputs: + - cuda_data + + - id: receiver + path: build/node + inputs: + cuda_data: sender/cuda_data + outputs: + - next diff --git a/examples/cmake-arrowcuda-dataflow/node/main.cc b/examples/cmake-arrowcuda-dataflow/node/main.cc new file mode 100644 index 0000000..4350613 --- /dev/null +++ b/examples/cmake-arrowcuda-dataflow/node/main.cc @@ -0,0 +1,178 @@ +#include "dora-node-api.h" +#include +#include +#include +#include +#include +#include + +// --------------------------------------------------------------------------- +// Arrow FFI helpers +// --------------------------------------------------------------------------- + +struct ArrowInput { + std::shared_ptr array; + std::string id; + rust::cxxbridge1::Box metadata; +}; + +ArrowInput receive_arrow(rust::cxxbridge1::Box event) { + struct ArrowArray c_array; + struct ArrowSchema c_schema; + + auto info = event_as_arrow_input_with_info( + std::move(event), + reinterpret_cast(&c_array), + reinterpret_cast(&c_schema)); + + if (!info.error.empty()) { + std::cerr << "receive_arrow error: " << std::string(info.error) + << std::endl; + return {nullptr, "", std::move(info.metadata)}; + } + + auto result = arrow::ImportArray(&c_array, &c_schema); + if (!result.ok()) { + std::cerr << "ImportArray failed: " << result.status().ToString() + << std::endl; + return {nullptr, "", std::move(info.metadata)}; + } + return {result.MoveValueUnsafe(), std::string(info.id), + std::move(info.metadata)}; +} + +bool send_arrow(rust::cxxbridge1::Box &sender, + const std::string &output_id, + const std::shared_ptr &array) { + struct ArrowArray c_array; + struct ArrowSchema c_schema; + + auto status = arrow::ExportArray(*array, &c_array, &c_schema); + if (!status.ok()) { + std::cerr << "ExportArray failed: " << status.ToString() << std::endl; + return false; + } + + auto result = send_arrow_output( + sender, output_id, + reinterpret_cast(&c_array), + reinterpret_cast(&c_schema)); + + if (!result.error.empty()) { + std::cerr << "send_arrow error: " << std::string(result.error) + << std::endl; + return false; + } + return true; +} + +// --------------------------------------------------------------------------- + +int main() { + auto dora_node = init_dora_node(); + auto id = std::string(node_id(dora_node.send_output)); + std::cout << "[" << id << "] started" << std::endl; + + // Get CUDA context for device 0 + auto manager = arrow::cuda::CudaDeviceManager::Instance(); + if (!manager.ok()) { + std::cerr << "Failed to get CudaDeviceManager: " + << manager.status().ToString() << std::endl; + return 1; + } + auto ctx_result = (*manager)->GetContext(0); + if (!ctx_result.ok()) { + std::cerr << "Failed to get CUDA context: " + << ctx_result.status().ToString() << std::endl; + return 1; + } + auto ctx = *ctx_result; + std::cout << "[" << id << "] CUDA context ready on device 0" << std::endl; + + // Build an empty array for "next" output + arrow::Int8Builder empty_builder; + auto empty_result = empty_builder.Finish(); + auto empty_array = empty_result.ValueOrDie(); + + for (int i = 0; i < 100; i++) { + auto event = dora_node.events->next(); + auto ty = event_type(event); + + if (ty == DoraEventType::Stop || + ty == DoraEventType::AllInputsClosed) { + break; + } + if (ty != DoraEventType::Input) { + continue; + } + + auto input = receive_arrow(std::move(event)); + if (!input.array) { + continue; + } + + // The Arrow array is Int8 containing the 64-byte CUDA IPC handle + auto int8_array = + std::static_pointer_cast(input.array); + const int8_t *raw = int8_array->raw_values(); + int64_t handle_len = int8_array->length(); + + // Read metadata + int64_t size = input.metadata->get_int("size"); + auto shape = input.metadata->get_list_int("shape"); + auto dtype_str = std::string(input.metadata->get_str("dtype")); + + std::cout << "[" << id << "] received CUDA IPC handle (" << handle_len + << " bytes), buffer size=" << size + << ", dtype=" << dtype_str << ", shape=["; + for (size_t j = 0; j < shape.size(); j++) { + if (j > 0) std::cout << ", "; + std::cout << shape[j]; + } + std::cout << "]" << std::endl; + + // Reconstruct IPC handle from raw bytes + auto ipc_handle_result = arrow::cuda::CudaIpcMemHandle::FromBuffer( + reinterpret_cast(raw)); + if (!ipc_handle_result.ok()) { + std::cerr << "FromBuffer failed: " + << ipc_handle_result.status().ToString() << std::endl; + continue; + } + auto ipc_handle = *ipc_handle_result; + + // Open the IPC buffer — zero-copy access to sender's GPU memory + auto buffer_result = ctx->OpenIpcBuffer(*ipc_handle); + if (!buffer_result.ok()) { + std::cerr << "OpenIpcBuffer failed: " + << buffer_result.status().ToString() << std::endl; + continue; + } + auto cuda_buffer = *buffer_result; + + // Copy first few values to host and print + int64_t copy_bytes = std::min(size, (int64_t)(10 * sizeof(int64_t))); + std::vector host_buf(copy_bytes); + auto copy_status = + cuda_buffer->CopyToHost(0, copy_bytes, host_buf.data()); + if (!copy_status.ok()) { + std::cerr << "CopyToHost failed: " << copy_status.ToString() + << std::endl; + continue; + } + + auto *values = reinterpret_cast(host_buf.data()); + int num_values = copy_bytes / sizeof(int64_t); + std::cout << "[" << id << "] first " << num_values << " values:"; + for (int j = 0; j < num_values; j++) { + std::cout << " " << values[j]; + } + std::cout << std::endl; + + // Send "next" to trigger the sender again + send_arrow(dora_node.send_output, "next", empty_array); + } + + std::cout << "[" << id << "] done" << std::endl; + return 0; +} diff --git a/examples/cmake-arrowcuda-dataflow/sender.py b/examples/cmake-arrowcuda-dataflow/sender.py new file mode 100644 index 0000000..3ce0a51 --- /dev/null +++ b/examples/cmake-arrowcuda-dataflow/sender.py @@ -0,0 +1,19 @@ +import pyarrow as pa +import torch +from dora import Node +from dora.cuda import torch_to_ipc_buffer + +torch.tensor([], device="cuda") + +node = Node() + +for i in range(10): + tensor = torch.arange(i * 100, (i + 1) * 100, dtype=torch.int64, device="cuda") + ipc_buffer, metadata = torch_to_ipc_buffer(tensor) + node.send_output("cuda_data", ipc_buffer, metadata) + + event = node.next() + if event["type"] != "INPUT": + break + +print("sender done") diff --git a/examples/rust-cuda-dataflow/dataflow.yml b/examples/rust-cuda-dataflow/dataflow.yml new file mode 100644 index 0000000..2c4b5cf --- /dev/null +++ b/examples/rust-cuda-dataflow/dataflow.yml @@ -0,0 +1,15 @@ +nodes: + - id: sender + path: sender.py + inputs: + next: receiver/next + outputs: + - cuda_data + + - id: receiver + build: cargo build -p rust-cuda-dataflow-receiver + path: ../../../target/debug/rust-cuda-dataflow-receiver + inputs: + cuda_data: sender/cuda_data + outputs: + - next diff --git a/examples/rust-cuda-dataflow/receiver/Cargo.toml b/examples/rust-cuda-dataflow/receiver/Cargo.toml new file mode 100644 index 0000000..9d24a08 --- /dev/null +++ b/examples/rust-cuda-dataflow/receiver/Cargo.toml @@ -0,0 +1,10 @@ +[package] +name = "rust-cuda-dataflow-receiver" +version = "0.1.0" +edition = "2024" +publish = false + +[dependencies] +dora-node-api = { git = "https://github.com/dora-rs/dora.git", rev = "77c277910b0ce87b902faa1ab369a33cbcd555f4" } +arrow-rs-cuda = { git = "https://github.com/haixuantao/arrow-rs-cuda" } +eyre = "0.6.8" diff --git a/examples/rust-cuda-dataflow/receiver/src/main.rs b/examples/rust-cuda-dataflow/receiver/src/main.rs new file mode 100644 index 0000000..766145f --- /dev/null +++ b/examples/rust-cuda-dataflow/receiver/src/main.rs @@ -0,0 +1,95 @@ +use arrow_rs_cuda::{CudaDeviceManager, CudaIpcMemHandle}; +use dora_node_api::{ + arrow::array::{Array, Int8Array}, + DoraNode, Event, MetadataParameters, Parameter, +}; +use eyre::Context; + +fn main() -> eyre::Result<()> { + let (mut node, mut events) = DoraNode::init_from_env()?; + + let manager = CudaDeviceManager::instance()?; + let ctx = manager.get_context(0)?; + println!("CUDA context ready on device 0"); + + while let Some(event) = events.recv() { + match event { + Event::Input { + id, + metadata, + data, + } => { + if id.as_str() != "cuda_data" { + continue; + } + + // Extract the 64-byte IPC handle from the Arrow Int8 array + let int8_array = data + .as_any() + .downcast_ref::() + .ok_or_else(|| eyre::eyre!("expected Int8Array"))?; + let values = int8_array.values(); + let handle_bytes = unsafe { + std::slice::from_raw_parts(values.as_ptr() as *const u8, values.len()) + }; + + // Read metadata + let size = match metadata.parameters.get("size") { + Some(Parameter::Integer(v)) => *v, + _ => eyre::bail!("missing or invalid 'size' metadata"), + }; + let shape = match metadata.parameters.get("shape") { + Some(Parameter::ListInt(v)) => v.clone(), + _ => eyre::bail!("missing or invalid 'shape' metadata"), + }; + let dtype = match metadata.parameters.get("dtype") { + Some(Parameter::String(v)) => v.clone(), + _ => eyre::bail!("missing or invalid 'dtype' metadata"), + }; + + println!( + "received CUDA IPC handle ({} bytes), buffer size={}, dtype={}, shape={:?}", + handle_bytes.len(), + size, + dtype, + shape, + ); + + // Open the IPC handle to get zero-copy access to sender's GPU memory + let ipc_handle = CudaIpcMemHandle::from_buffer(handle_bytes) + .wrap_err("failed to create IPC handle from buffer")?; + let cuda_buffer = ctx + .open_ipc_buffer(&ipc_handle) + .wrap_err("failed to open IPC buffer")?; + + // Copy first 10 int64 values to host and print + let copy_bytes = std::cmp::min(size, 10 * 8) as usize; + let mut host_buf = vec![0u8; copy_bytes]; + cuda_buffer + .copy_to_host(0, &mut host_buf) + .wrap_err("failed to copy to host")?; + + let num_values = copy_bytes / std::mem::size_of::(); + let values = + unsafe { std::slice::from_raw_parts(host_buf.as_ptr() as *const i64, num_values) }; + print!("first {} values:", num_values); + for v in values { + print!(" {}", v); + } + println!(); + + // Send "next" to trigger the sender again + node.send_output( + "next".into(), + MetadataParameters::default(), + dora_node_api::arrow::array::UInt8Array::from(vec![0u8]), + )?; + } + Event::Stop(_) => break, + _ => {} + } + } + + println!("receiver done"); + Ok(()) +} diff --git a/examples/rust-cuda-dataflow/sender.py b/examples/rust-cuda-dataflow/sender.py new file mode 100644 index 0000000..3ce0a51 --- /dev/null +++ b/examples/rust-cuda-dataflow/sender.py @@ -0,0 +1,19 @@ +import pyarrow as pa +import torch +from dora import Node +from dora.cuda import torch_to_ipc_buffer + +torch.tensor([], device="cuda") + +node = Node() + +for i in range(10): + tensor = torch.arange(i * 100, (i + 1) * 100, dtype=torch.int64, device="cuda") + ipc_buffer, metadata = torch_to_ipc_buffer(tensor) + node.send_output("cuda_data", ipc_buffer, metadata) + + event = node.next() + if event["type"] != "INPUT": + break + +print("sender done")