PyO3 Advanced Type Conversion Reference
Version: 1.0.0 Last Updated: 2025-10-30 PyO3 Version: 0.20+ Python Version: 3.8+ Rust Version: 1.70+ rust-numpy: 0.20+ ndarray: 0.15+ arrow: 48.0+
Table of Contents
- Introduction
- Zero-Copy Fundamentals
- Numpy Integration
- Buffer Protocol
- Arrow/Parquet Integration
- Custom Conversion Protocols
- Streaming Conversion
- Performance Optimization
- Memory Safety Patterns
- Best Practices
- Troubleshooting
Introduction
Overview
Advanced type conversion in PyO3 focuses on high-performance data interchange between Rust and Python with minimal overhead. This includes zero-copy operations, numpy array integration, Apache Arrow/Parquet support, and custom protocol implementations.
Key Concepts
- Zero-Copy: Share memory between languages without duplication
- Numpy Integration: Efficient ndarray operations via rust-numpy
- Buffer Protocol: Python's low-level memory interface
- Arrow/Parquet: Columnar data formats for analytics
- Streaming: Process large datasets incrementally
Performance Benefits
| Operation | Copy | Zero-Copy | Speedup |
|---|---|---|---|
| 1GB numpy array pass | ~1000ms | ~1ms | 1000x |
| 100MB Arrow table | ~200ms | ~1ms | 200x |
| Large dataset (10GB) | OOM | Streaming | ∞ |
Zero-Copy Fundamentals
Memory Sharing Basics
use pyo3::prelude::*;
use pyo3::types::PyBytes;
// Copy: Data duplicated
#[pyfunction]
fn process_copy(data: Vec<u8>) -> Vec<u8> {
// data is copied from Python to Rust
let processed: Vec<u8> = data.iter().map(|x| x.wrapping_add(1)).collect();
// result is copied from Rust to Python
processed
}
// Zero-copy read: Borrow Python data
#[pyfunction]
fn process_zerocopy_read(data: &[u8]) -> Vec<u8> {
// data is borrowed (no copy)
data.iter().map(|x| x.wrapping_add(1)).collect()
// result still copied (return value)
}
// Zero-copy write: Return view into Rust data
#[pyfunction]
fn process_zerocopy_write(py: Python, size: usize) -> PyObject {
let data: Vec<u8> = (0..size).map(|i| (i % 256) as u8).collect();
// Transfer ownership to Python (no copy)
PyBytes::new(py, &data).into()
}
Python usage:
# Copy: 2 copies (Python→Rust, Rust→Python)
data = bytes(1_000_000)
result = process_copy(data)
# Zero-copy read: 1 copy (Rust→Python)
result = process_zerocopy_read(data)
# Zero-copy write: 0 copies
result = process_zerocopy_write(1_000_000)
Safety Guarantees
// UNSAFE: Dangling reference
#[pyfunction]
fn unsafe_reference<'a>(data: &'a [u8]) -> &'a [u8] {
&data[..10] // Returns reference to borrowed data
// DANGER: Reference may outlive borrowed data
}
// SAFE: Own the data
#[pyfunction]
fn safe_slice(data: &[u8]) -> Vec<u8> {
data[..10].to_vec() // Copy the slice (owns data)
}
// SAFE: Use PyO3 lifetime management
#[pyfunction]
fn safe_bytes(py: Python, data: &[u8]) -> PyObject {
PyBytes::new(py, &data[..10]).into() // PyO3 manages lifetime
}
Lifetime Management
use pyo3::types::PyByteArray;
#[pyfunction]
fn modify_in_place(array: &PyByteArray) -> PyResult<()> {
// Safe: PyByteArray provides mutable access
let mut bytes = unsafe { array.as_bytes_mut() };
for b in bytes.iter_mut() {
*b = b.wrapping_add(1);
}
Ok(())
}
Python usage:
data = bytearray(b"hello")
modify_in_place(data)
print(data) # bytearray(b'ifmmp')
Numpy Integration
rust-numpy Basics
use numpy::{PyArray1, PyArray2, PyReadonlyArray1, PyReadonlyArray2};
use ndarray::{Array1, Array2};
// Read-only array access (zero-copy read)
#[pyfunction]
fn sum_array(array: PyReadonlyArray1<f64>) -> f64 {
let array = array.as_array(); // ndarray view (no copy)
array.sum()
}
// Mutable array access
#[pyfunction]
fn double_array(py: Python, mut array: PyArray1<f64>) {
let array = unsafe { array.as_array_mut() }; // Mutable view
array.mapv_inplace(|x| x * 2.0);
}
// Create new array (zero-copy to Python)
#[pyfunction]
fn create_array(py: Python, size: usize) -> &PyArray1<f64> {
let data: Vec<f64> = (0..size).map(|i| i as f64).collect();
PyArray1::from_vec(py, data) // Transfers ownership to Python
}
Python usage:
import numpy as np
arr = np.array([1.0, 2.0, 3.0, 4.0, 5.0])
total = sum_array(arr) # Zero-copy read
arr = np.array([1.0, 2.0, 3.0])
double_array(arr) # In-place modification
print(arr) # [2.0, 4.0, 6.0]
arr = create_array(1_000_000) # Zero-copy creation
Multi-dimensional Arrays
#[pyfunction]
fn transpose_matrix(array: PyReadonlyArray2<f64>) -> Py<PyArray2<f64>> {
let array = array.as_array();
let transposed = array.t().to_owned(); // Transpose
Python::with_gil(|py| {
PyArray2::from_array(py, &transposed).into()
})
}
#[pyfunction]
fn matrix_multiply(
a: PyReadonlyArray2<f64>,
b: PyReadonlyArray2<f64>,
) -> Py<PyArray2<f64>> {
let a = a.as_array();
let b = b.as_array();
let result = a.dot(&b);
Python::with_gil(|py| {
PyArray2::from_array(py, &result).into()
})
}
Custom dtypes
use numpy::PyArray;
#[repr(C)]
#[derive(Clone, Copy)]
struct Point {
x: f64,
y: f64,
z: f64,
}
unsafe impl numpy::Element for Point {
const IS_COPY: bool = true;
fn get_dtype(py: Python) -> &numpy::PyArrayDescr {
numpy::PyArrayDescr::new(
py,
&[
("x", numpy::npyffi::NPY_FLOAT64),
("y", numpy::npyffi::NPY_FLOAT64),
("z", numpy::npyffi::NPY_FLOAT64),
],
)
}
}
#[pyfunction]
fn process_points(points: PyReadonlyArray1<Point>) -> f64 {
let points = points.as_array();
points.iter().map(|p| (p.x.powi(2) + p.y.powi(2) + p.z.powi(2)).sqrt()).sum()
}
Python usage:
import numpy as np
# Create structured array
dtype = np.dtype([('x', 'f8'), ('y', 'f8'), ('z', 'f8')])
points = np.array([(1.0, 2.0, 3.0), (4.0, 5.0, 6.0)], dtype=dtype)
total_distance = process_points(points)
Buffer Protocol
Implementing Buffer Protocol
use pyo3::buffer::PyBuffer;
use pyo3::ffi;
#[pyclass]
struct CustomBuffer {
data: Vec<u8>,
}
#[pymethods]
impl CustomBuffer {
#[new]
fn new(size: usize) -> Self {
CustomBuffer {
data: vec![0u8; size],
}
}
fn __getbuffer__(&self, view: *mut ffi::Py_buffer, flags: std::os::raw::c_int) -> PyResult<()> {
// Implement buffer protocol
unsafe {
(*view).buf = self.data.as_ptr() as *mut std::os::raw::c_void;
(*view).len = self.data.len() as isize;
(*view).readonly = 0;
(*view).itemsize = 1;
(*view).format = b"B".as_ptr() as *mut std::os::raw::c_char;
(*view).ndim = 1;
(*view).shape = std::ptr::null_mut();
(*view).strides = std::ptr::null_mut();
(*view).suboffsets = std::ptr::null_mut();
(*view).internal = std::ptr::null_mut();
}
Ok(())
}
fn __releasebuffer__(&self, _view: *mut ffi::Py_buffer) {
// Cleanup if needed
}
}
Python usage:
buf = CustomBuffer(1000)
mv = memoryview(buf) # Access via buffer protocol
print(mv.nbytes) # 1000
print(mv.format) # 'B' (unsigned byte)
# Use with numpy
arr = np.array(buf, copy=False) # Zero-copy view
Reading Buffers
#[pyfunction]
fn process_buffer(py: Python, obj: &PyAny) -> PyResult<usize> {
let buffer = PyBuffer::get(py, obj)?;
// Check buffer properties
println!("Buffer info:");
println!(" Length: {}", buffer.len_bytes());
println!(" Readonly: {}", buffer.readonly());
println!(" Dimensions: {}", buffer.dimensions());
// Access data
let slice = unsafe { buffer.as_slice::<u8>(py)? };
Ok(slice.iter().map(|&x| x as usize).sum())
}
Python usage:
data = b"hello world"
total = process_buffer(data)
arr = np.array([1, 2, 3, 4, 5], dtype=np.uint8)
total = process_buffer(arr)
Arrow/Parquet Integration
Arrow Arrays
use arrow::array::{Array, Int64Array, Float64Array};
use arrow::datatypes::{Schema, Field, DataType};
use arrow::record_batch::RecordBatch;
use pyo3::types::PyList;
#[pyfunction]
fn create_arrow_array(py: Python, data: Vec<i64>) -> PyResult<PyObject> {
let array = Int64Array::from(data);
// Convert to Python (using pyarrow)
let pyarrow = py.import("pyarrow")?;
// This is simplified - actual implementation uses FFI
let py_array = pyarrow.call_method1("array", (array.values().as_slice(),))?;
Ok(py_array.into())
}
#[pyfunction]
fn process_arrow_batch(py: Python, batch: &PyAny) -> PyResult<f64> {
// Import pyarrow
let pyarrow = py.import("pyarrow")?;
// Get column
let column = batch.call_method1("column", (0,))?;
// Convert to Rust (zero-copy via Arrow C Data Interface)
let array: &Int64Array = /* conversion via FFI */;
// Process in Rust
let sum: i64 = array.values().iter().sum();
Ok(sum as f64)
}
Parquet Files
use arrow::record_batch::RecordBatch;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::file::reader::FileReader;
#[pyfunction]
fn read_parquet(py: Python, path: &str) -> PyResult<PyObject> {
use std::fs::File;
let file = File::open(path)
.map_err(|e| PyIOError::new_err(e.to_string()))?;
let builder = ParquetRecordBatchReaderBuilder::try_new(file)
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
let reader = builder.build()
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
// Read all batches
let mut batches = Vec::new();
for batch_result in reader {
let batch = batch_result
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
batches.push(batch);
}
// Convert to Python pyarrow Table
let pyarrow = py.import("pyarrow")?;
// Conversion implementation...
Ok(py.None())
}
#[pyfunction]
fn write_parquet(path: &str, data: Vec<Vec<i64>>) -> PyResult<()> {
use parquet::arrow::arrow_writer::ArrowWriter;
use std::fs::File;
use std::sync::Arc;
// Create Arrow schema
let schema = Schema::new(vec![
Field::new("column1", DataType::Int64, false),
Field::new("column2", DataType::Int64, false),
]);
// Create RecordBatch
let arrays: Vec<Arc<dyn Array>> = vec![
Arc::new(Int64Array::from(data[0].clone())),
Arc::new(Int64Array::from(data[1].clone())),
];
let batch = RecordBatch::try_new(Arc::new(schema.clone()), arrays)
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
// Write to Parquet
let file = File::create(path)
.map_err(|e| PyIOError::new_err(e.to_string()))?;
let mut writer = ArrowWriter::try_new(file, Arc::new(schema), None)
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
writer.write(&batch)
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
writer.close()
.map_err(|e| PyRuntimeError::new_err(e.to_string()))?;
Ok(())
}
Custom Conversion Protocols
Sequence Protocol
#[pyclass]
struct RustVec {
data: Vec<i64>,
}
#[pymethods]
impl RustVec {
#[new]
fn new() -> Self {
RustVec { data: Vec::new() }
}
fn __len__(&self) -> usize {
self.data.len()
}
fn __getitem__(&self, index: isize) -> PyResult<i64> {
let idx = if index < 0 {
(self.data.len() as isize + index) as usize
} else {
index as usize
};
self.data.get(idx).copied()
.ok_or_else(|| PyIndexError::new_err("Index out of range"))
}
fn __setitem__(&mut self, index: isize, value: i64) -> PyResult<()> {
let idx = if index < 0 {
(self.data.len() as isize + index) as usize
} else {
index as usize
};
if idx < self.data.len() {
self.data[idx] = value;
Ok(())
} else {
Err(PyIndexError::new_err("Index out of range"))
}
}
fn append(&mut self, value: i64) {
self.data.push(value);
}
}
Mapping Protocol
use std::collections::HashMap;
#[pyclass]
struct RustDict {
data: HashMap<String, i64>,
}
#[pymethods]
impl RustDict {
#[new]
fn new() -> Self {
RustDict { data: HashMap::new() }
}
fn __len__(&self) -> usize {
self.data.len()
}
fn __getitem__(&self, key: String) -> PyResult<i64> {
self.data.get(&key).copied()
.ok_or_else(|| PyKeyError::new_err(format!("Key not found: {}", key)))
}
fn __setitem__(&mut self, key: String, value: i64) {
self.data.insert(key, value);
}
fn __delitem__(&mut self, key: String) -> PyResult<()> {
self.data.remove(&key)
.ok_or_else(|| PyKeyError::new_err(format!("Key not found: {}", key)))?;
Ok(())
}
fn __contains__(&self, key: String) -> bool {
self.data.contains_key(&key)
}
fn keys(&self) -> Vec<String> {
self.data.keys().cloned().collect()
}
fn values(&self) -> Vec<i64> {
self.data.values().copied().collect()
}
fn items(&self) -> Vec<(String, i64)> {
self.data.iter().map(|(k, v)| (k.clone(), *v)).collect()
}
}
Streaming Conversion
Chunked Processing
use std::fs::File;
use std::io::{BufRead, BufReader};
#[pyclass]
struct ChunkReader {
reader: Option<BufReader<File>>,
chunk_size: usize,
}
#[pymethods]
impl ChunkReader {
#[new]
fn new(path: String, chunk_size: usize) -> PyResult<Self> {
let file = File::open(&path)
.map_err(|e| PyIOError::new_err(e.to_string()))?;
Ok(ChunkReader {
reader: Some(BufReader::new(file)),
chunk_size,
})
}
fn __iter__(slf: PyRef<Self>) -> PyRef<Self> {
slf
}
fn __next__(mut slf: PyRefMut<Self>) -> PyResult<Option<Vec<String>>> {
if let Some(ref mut reader) = slf.reader {
let mut chunk = Vec::new();
for _ in 0..slf.chunk_size {
let mut line = String::new();
match reader.read_line(&mut line) {
Ok(0) => break, // EOF
Ok(_) => chunk.push(line.trim_end().to_string()),
Err(e) => return Err(PyIOError::new_err(e.to_string())),
}
}
if chunk.is_empty() {
Ok(None)
} else {
Ok(Some(chunk))
}
} else {
Ok(None)
}
}
}
Python usage:
reader = ChunkReader("large_file.txt", chunk_size=1000)
for chunk in reader:
process_chunk(chunk) # Process 1000 lines at a time
Incremental Conversion
#[pyclass]
struct IncrementalProcessor {
buffer: Vec<u8>,
processed: usize,
}
#[pymethods]
impl IncrementalProcessor {
#[new]
fn new(size: usize) -> Self {
IncrementalProcessor {
buffer: vec![0u8; size],
processed: 0,
}
}
fn process_chunk(&mut self, size: usize) -> Option<Vec<u8>> {
if self.processed >= self.buffer.len() {
return None;
}
let end = (self.processed + size).min(self.buffer.len());
let chunk = self.buffer[self.processed..end].to_vec();
self.processed = end;
Some(chunk.iter().map(|x| x.wrapping_add(1)).collect())
}
fn progress(&self) -> f64 {
self.processed as f64 / self.buffer.len() as f64
}
}
Performance Optimization
SIMD Optimization
#[cfg(target_arch = "x86_64")]
use std::arch::x86_64::*;
#[pyfunction]
fn simd_sum(data: &[f64]) -> f64 {
#[cfg(target_arch = "x86_64")]
{
if is_x86_feature_detected!("avx2") {
return unsafe { simd_sum_avx2(data) };
}
}
// Fallback
data.iter().sum()
}
#[cfg(target_arch = "x86_64")]
#[target_feature(enable = "avx2")]
unsafe fn simd_sum_avx2(data: &[f64]) -> f64 {
let mut sum = _mm256_setzero_pd();
let chunks = data.chunks_exact(4);
let remainder = chunks.remainder();
for chunk in chunks {
let vals = _mm256_loadu_pd(chunk.as_ptr());
sum = _mm256_add_pd(sum, vals);
}
// Horizontal sum
let mut result = [0.0f64; 4];
_mm256_storeu_pd(result.as_mut_ptr(), sum);
result.iter().sum::<f64>() + remainder.iter().sum::<f64>()
}
Cache-Friendly Layouts
// Bad: Structure of Arrays (cache-unfriendly for iteration)
struct ParticlesSOA {
x: Vec<f64>,
y: Vec<f64>,
z: Vec<f64>,
mass: Vec<f64>,
}
// Good: Array of Structures (cache-friendly)
#[repr(C)]
struct Particle {
x: f64,
y: f64,
z: f64,
mass: f64,
}
struct ParticlesAOS {
particles: Vec<Particle>,
}
#[pyfunction]
fn compute_energy(particles: &[Particle]) -> f64 {
// All fields accessed together - better cache locality
particles.iter()
.map(|p| 0.5 * p.mass * (p.x.powi(2) + p.y.powi(2) + p.z.powi(2)))
.sum()
}
Parallel Conversion
use rayon::prelude::*;
#[pyfunction]
fn parallel_process(py: Python, data: Vec<f64>) -> Vec<f64> {
py.allow_threads(|| {
data.par_iter()
.map(|&x| expensive_computation(x))
.collect()
})
}
fn expensive_computation(x: f64) -> f64 {
// Expensive operation
(0..1000).fold(x, |acc, _| acc.sin().cos())
}
Memory Safety Patterns
Preventing Use-After-Free
// UNSAFE: Returning reference to temporary
#[pyfunction]
fn unsafe_slice<'a>() -> &'a [u8] {
let data = vec![1, 2, 3];
&data // DANGER: data dropped, reference dangles
}
// SAFE: Transfer ownership
#[pyfunction]
fn safe_vec(py: Python) -> PyObject {
let data = vec![1, 2, 3];
PyBytes::new(py, &data).into() // Python owns the data
}
Thread Safety
use std::sync::Arc;
use std::sync::Mutex;
#[pyclass]
struct ThreadSafeCounter {
value: Arc<Mutex<i64>>,
}
#[pymethods]
impl ThreadSafeCounter {
#[new]
fn new() -> Self {
ThreadSafeCounter {
value: Arc::new(Mutex::new(0)),
}
}
fn increment(&self) {
let mut value = self.value.lock().unwrap();
*value += 1;
}
fn get(&self) -> i64 {
*self.value.lock().unwrap()
}
}
Best Practices
1. Choose Appropriate Conversion Strategy
// Small data (< 1MB): Copy is fine
#[pyfunction]
fn process_small(data: Vec<u8>) -> Vec<u8> {
data.iter().map(|x| x.wrapping_add(1)).collect()
}
// Large data (> 1MB): Use zero-copy
#[pyfunction]
fn process_large(data: &[u8]) -> Vec<u8> {
data.iter().map(|x| x.wrapping_add(1)).collect()
}
// Very large data (> 100MB): Stream
#[pyfunction]
fn process_huge(py: Python, data: &[u8], chunk_size: usize) -> PyObject {
// Return iterator for chunked processing
// Implementation...
py.None()
}
2. Validate Array Properties
#[pyfunction]
fn process_array(array: PyReadonlyArray2<f64>) -> PyResult<f64> {
let array = array.as_array();
// Validate shape
let shape = array.shape();
if shape[0] == 0 || shape[1] == 0 {
return Err(PyValueError::new_err("Array cannot be empty"));
}
// Validate contiguous
if !array.is_standard_layout() {
return Err(PyValueError::new_err("Array must be C-contiguous"));
}
Ok(array.sum())
}
3. Profile Before Optimizing
import time
import numpy as np
def benchmark(func, *args, iterations=100):
start = time.time()
for _ in range(iterations):
func(*args)
end = time.time()
return (end - start) / iterations
data = np.random.rand(1_000_000)
copy_time = benchmark(process_copy, data)
zerocopy_time = benchmark(process_zerocopy, data)
print(f"Copy: {copy_time*1000:.2f}ms")
print(f"Zero-copy: {zerocopy_time*1000:.2f}ms")
print(f"Speedup: {copy_time/zerocopy_time:.1f}x")
Troubleshooting
Common Issues
1. Array not contiguous:
# Problem: Non-contiguous array
arr = np.array([[1, 2], [3, 4]])[:, 0] # Non-contiguous
# Solution: Make contiguous
arr = np.ascontiguousarray(arr)
result = process_array(arr)
2. Dtype mismatch:
# Problem: Wrong dtype
arr = np.array([1, 2, 3], dtype=np.int32) # Expecting f64
# Solution: Convert dtype
arr = arr.astype(np.float64)
result = process_array(arr)
3. Memory not released:
// Problem: Holding reference too long
let array = array.as_array(); // Borrows array
// Do work...
drop(array); // Explicitly release
// Or use scope
{
let array = array.as_array();
// Work here
} // Automatically released
Conclusion
This reference covered:
- Zero-copy fundamentals: Memory sharing, safety, lifetimes
- Numpy integration: rust-numpy, multi-dimensional arrays, custom dtypes
- Buffer protocol: Implementation, reading buffers
- Arrow/Parquet: Arrays, record batches, file I/O
- Custom protocols: Sequence, mapping, iteration
- Streaming: Chunked processing, incremental conversion
- Performance: SIMD, cache optimization, parallelization
- Memory safety: Preventing errors, thread safety
- Best practices: Strategy selection, validation, profiling
Document Version: 1.0.0 Lines: 1,100+ Last Updated: 2025-10-30