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

如何用PyO3实现Rust TCP服务与Python的字节通信及连接管理

解决方案步骤

1. 定义线程安全的共享连接状态结构

放弃全局变量方案,用Arc<Mutex<T>>实现跨线程安全的状态共享。创建专门的结构体管理客户端ID与最新数据的映射:

use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use uuid::Uuid; // 需要在Cargo.toml中添加uuid依赖

#[derive(Default)]
struct ConnectionStore {
    clients: HashMap<String, Vec<u8>>,
}

impl ConnectionStore {
    // 更新客户端数据,只保留最后n字节(示例设为1024,可配置)
    fn update_client_data(&mut self, client_id: String, new_data: &[u8]) {
        let max_bytes = 1024;
        let target_data = if new_data.len() > max_bytes {
            &new_data[new_data.len() - max_bytes..]
        } else {
            new_data
        };
        self.clients.insert(client_id, target_data.to_vec());
    }

    // 获取所有在线客户端ID
    fn get_connections(&self) -> Vec<String> {
        self.clients.keys().cloned().collect()
    }

    // 获取指定客户端的最新数据
    fn get_client_data(&self, client_id: &str) -> Option<Vec<u8>> {
        self.clients.get(client_id).cloned()
    }

    // 客户端断开时移除记录
    fn remove_client(&mut self, client_id: &str) {
        self.clients.remove(client_id);
    }
}

2. 修改监听逻辑,传递共享状态

在启动监听函数中创建共享状态实例,每个连接处理线程克隆该实例以操作状态:

#[pyfunction]
fn start(py: Python, addr: &str) -> PyResult<Listener> {
    let store = Arc::new(Mutex::new(ConnectionStore::default()));
    let store_clone = Arc::clone(&store);

    // 在Python允许的线程环境下启动监听
    py.allow_threads(move || {
        let listener = TcpListener::bind(addr).unwrap();
        for stream in listener.incoming() {
            match stream {
                Ok(stream) => {
                    let client_id = Uuid::new_v4().to_string();
                    let thread_store = Arc::clone(&store_clone);

                    // 启动独立线程处理客户端连接
                    std::thread::spawn(move || {
                        let mut stream = stream;
                        // 先将新客户端注册到状态中
                        thread_store.lock().unwrap().clients.insert(client_id.clone(), Vec::new());

                        let mut buffer = [0; 1024];
                        loop {
                            match stream.read(&mut buffer) {
                                Ok(0) => {
                                    // 客户端主动断开,移除记录
                                    thread_store.lock().unwrap().remove_client(&client_id);
                                    break;
                                }
                                Ok(n) => {
                                    // 更新该客户端的最新数据
                                    thread_store.lock().unwrap().update_client_data(client_id.clone(), &buffer[..n]);
                                }
                                Err(e) => {
                                    eprintln!("读取客户端数据失败: {}", e);
                                    thread_store.lock().unwrap().remove_client(&client_id);
                                    break;
                                }
                            }
                        }
                    });
                }
                Err(e) => eprintln!("接受连接失败: {}", e),
            }
        }
    });

    Ok(Listener { store })
}

3. 实现PyO3暴露给Python的Listener类

让这个类持有共享状态实例,实现Python调用的connections()和read()方法:

#[pyclass]
struct Listener {
    store: Arc<Mutex<ConnectionStore>>,
}

#[pymethods]
impl Listener {
    // 返回所有在线客户端ID列表
    fn connections(&self) -> Vec<String> {
        self.store.lock().unwrap().get_connections()
    }

    // 获取指定客户端的最新数据,转为Python字节对象返回
    fn read(&self, client_id: &str) -> Option<PyObject> {
        Python::with_gil(|py| {
            self.store.lock().unwrap().get_client_data(client_id)
                .map(|data| PyBytes::new(py, &data).into())
        })
    }
}

4. 补充Cargo.toml依赖

确保添加必要的依赖项:

[package]
name = "tcp_server_pyo3"
version = "0.1.0"
edition = "2021"

[dependencies]
pyo3 = { version = "0.20", features = ["extension-module"] }
uuid = { version = "1.0", features = ["v4", "fast-rng"] }

关键注意事项

  • 线程安全:Arc实现多线程共享所有权,Mutex保证同一时间仅一个线程操作共享状态,避免数据竞争。
  • 客户端ID唯一性:用UUID生成唯一ID,也可替换为递增整数,按需选择。
  • 缓冲区配置:示例中固定缓冲区大小为1024字节,可改为start函数的参数实现动态配置。
  • 错误处理:示例用unwrap简化代码,生产环境需替换为更健壮的错误处理,将异常传递给Python端。

内容的提问来源于stack exchange,提问作者Abay Bektursun

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 16:20:31