如何用Rust持续采集传感器时序数据并供Python查询(无内存复制)
Great question! Let’s walk through how to build this system step by step, and tackle that zero-copy data access question you’re curious about—spoiler: it’s absolutely possible.
Here’s the high-level plan:
- Rust layer: Handles high-performance sensor data collection, stores time-series data in a thread-safe, efficient buffer (like a ring buffer), and exposes query functions via a Python-compatible interface.
- Python layer: Calls these Rust functions to fetch filtered data (like the last 5 minutes of readings) and works with it using familiar tools like numpy.
1. Sensor Data Collection & Storage
First, define a lightweight struct to hold each sensor reading—critical for handling high-volume time-series data efficiently:
#[derive(Debug, Clone, Copy)] struct SensorReading { timestamp: u64, // Millisecond timestamp for precise time filtering value: f64, // The sensor's measurement }
Use an async runtime like Tokio to run a continuous collection loop. We’ll store readings in a thread-safe ring buffer (I recommend the ringbuf crate or crossbeam-channel for bounded, lock-free storage) wrapped in Arc<Mutex> so multiple threads can access it safely:
use tokio; use ringbuf::RingBuffer; use std::sync::{Arc, Mutex}; use chrono::Utc; #[tokio::main] async fn main() { let buffer = Arc::new(Mutex::new(RingBuffer::<SensorReading>::new(10000))); let collector_buffer = Arc::clone(&buffer); // Spawn a background task to collect sensor data tokio::spawn(async move { loop { // Replace this with your actual sensor reading logic let reading = SensorReading { timestamp: Utc::now().timestamp_millis() as u64, value: rand::random::<f64>() * 100.0, }; let mut buf = collector_buffer.lock().unwrap(); buf.push(reading).unwrap_or(()); // Drop oldest reading if buffer is full tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; } }); // Keep the main thread alive (or integrate with your Python extension setup here) tokio::signal::ctrl_c().await.unwrap(); }
2. Expose Query Functions to Python
Use the PyO3 crate to wrap our Rust query logic into a Python extension. This lets Python call Rust functions natively. Here’s how to implement a function that fetches the last 5 minutes of readings:
use pyo3::prelude::*; use pyo3::types::PyArray2; use chrono::Utc; use std::sync::{Arc, Mutex}; use ringbuf::RingBuffer; use once_cell::sync::Lazy; // Shared buffer, initialized once static SHARED_BUFFER: Lazy<Arc<Mutex<RingBuffer<SensorReading>>>> = Lazy::new(|| { Arc::new(Mutex::new(RingBuffer::<SensorReading>::new(10000))) }); #[pyfunction] fn get_last_5_minutes(py: Python) -> PyResult<&PyArray2<f64>> { let buffer = SHARED_BUFFER.lock().unwrap(); let now = Utc::now().timestamp_millis() as u64; let five_min_ago = now - (5 * 60 * 1000); // Filter readings from the last 5 minutes let filtered: Vec<SensorReading> = buffer .iter() .filter(|r| r.timestamp >= five_min_ago) .cloned() .collect(); // Convert to a flat Vec<f64> for numpy: [timestamp1, value1, timestamp2, value2, ...] let mut data = Vec::with_capacity(filtered.len() * 2); for r in filtered { data.push(r.timestamp as f64); data.push(r.value); } // Reshape into a 2D numpy array (N rows, 2 columns: timestamp + value) let array = PyArray2::from_slice(py, &data)?.reshape((filtered.len(), 2))?; Ok(array) } #[pymodule] fn sensor_collector(_py: Python, m: &PyModule) -> PyResult<()> { m.add_function(wrap_pyfunction!(get_last_5_minutes, m)?)?; Ok(()) }
Once you compile the Rust extension (follow PyO3’s setup docs), you can call it just like any Python module:
import sensor_collector import numpy as np # Fetch last 5 minutes of data readings = sensor_collector.get_last_5_minutes() print(f"Got {readings.shape[0]} readings") print(readings[:5]) # Print first 5 timestamp-value pairs
Short answer: Yes, but it requires careful handling of memory ownership. Here’s why the default approach copies data, and how to fix it:
Why Copy Happens By Default
In the example above, we convert a Rust Vec into a numpy array. PyO3 copies the data because Rust’s memory is managed via ownership rules, while Python uses garbage collection. If we just passed a pointer to Rust’s Vec, Python wouldn’t know when it’s safe to free that memory—leading to dangling pointers or crashes.
How to Implement Zero-Copy
To let Python access Rust’s memory directly without copying, we need to:
- Ensure the Rust memory is continuous (ring buffers can wrap around, so you might need to copy only if valid data is split across the buffer’s head and tail).
- Keep Rust’s memory alive for as long as Python is using it. We can do this by attaching an
Arcreference to the numpy array, so Python’s GC knows not to free the memory until the array is discarded.
Here’s a modified zero-copy version of the query function:
use pyo3::types::PyArray2; use pyo3::prelude::*; use std::sync::Arc; #[pyfunction] fn get_last_5_minutes_zero_copy(py: Python) -> PyResult<Py<PyArray2<f64>>> { let buffer = SHARED_BUFFER.lock().unwrap(); let now = Utc::now().timestamp_millis() as u64; let five_min_ago = now - (5 * 60 * 1000); // Get raw slices from the ring buffer let (head, tail) = buffer.as_slices(); let mut filtered_slice: &[SensorReading] = &[]; // Check if valid data is in one continuous slice or split if let Some(first_valid_idx) = head.iter().position(|r| r.timestamp >= five_min_ago) { filtered_slice = &head[first_valid_idx..]; } else if !tail.is_empty() && tail.iter().all(|r| r.timestamp >= five_min_ago) { filtered_slice = tail; } // Handle continuous vs split data let (data_ptr, len, arc_holder) = if filtered_slice.is_empty() { // Empty case: return empty array let empty_vec = Vec::new(); (empty_vec.as_ptr(), 0, Arc::clone(&SHARED_BUFFER)) } else if filtered_slice.as_ptr() == head.as_ptr() || filtered_slice.as_ptr() == tail.as_ptr() { // Continuous slice: use raw pointer directly (filtered_slice.as_ptr() as *const f64, filtered_slice.len(), Arc::clone(&SHARED_BUFFER)) } else { // Split data: copy to continuous Vec (unavoidable here) let filtered: Vec<SensorReading> = head.iter().chain(tail.iter()) .filter(|r| r.timestamp >= five_min_ago) .cloned() .collect(); (filtered.as_ptr() as *const f64, filtered.len(), Arc::new(Mutex::new(RingBuffer::from_vec(filtered)))) }; // Create numpy array from raw pointer let dims = [len, 2]; let array = unsafe { PyArray2::from_raw_parts(py, dims, data_ptr) }?; // Attach Arc to numpy array so Python GC keeps memory alive let arc_py = Py::new(py, arc_holder)?; array.set_user_data(py, arc_py)?; Ok(array.into()) }
Alternative: Use Apache Arrow for Zero-Copy
For a more robust, production-ready zero-copy solution, use Apache Arrow. Arrow defines a cross-language memory format that both Rust and Python can read directly:
- In Rust, use the
arrowcrate to create Arrow arrays from your sensor data. - In Python, use
pyarrowto receive these arrays and convert them to numpy arrays (no copy needed—numpy can view Arrow’s memory directly).
This avoids manual unsafe pointer handling and works seamlessly across languages.
- Rust is ideal for sensor data collection because it’s fast, low-overhead, and safe for concurrent code.
- Zero-copy is achievable, but depends on whether your data is stored in continuous memory. For ring buffers, you’ll need to handle split slices (copying only when necessary).
- Apache Arrow is the easiest way to get reliable zero-copy between Rust and Python without manual memory management.
内容的提问来源于stack exchange,提问作者Greg

