You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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.

整体 System Architecture

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.
Rust Side Implementation

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(())
}
Python Side Usage

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
Zero-Copy Data Access: Is It Possible?

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:

  1. 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).
  2. Keep Rust’s memory alive for as long as Python is using it. We can do this by attaching an Arc reference 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 arrow crate to create Arrow arrays from your sensor data.
  • In Python, use pyarrow to 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.

Final Notes
  • 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 08:29:15