如何创建可更新共享数据的多Future且满足spawn trait约束?
我正尝试构建一个客户端/服务器框架,让多个客户端向管理共享数据的单服务器发送请求。以下是简化版代码:
use async_std::task; use futures::channel::mpsc::{self, Receiver}; use futures::{SinkExt, StreamExt}; use std::sync::{Arc, Mutex}; struct Service { handlers: Vec<Receiver<bool>>, data: Arc<Mutex<bool>>, } impl Service { async fn handle_requests(&mut self) { let mut handlers = Vec::new(); for rx in &mut self.handlers { let data_arc = self.data.clone(); handlers.push(Box::pin(async move { let mut data = data_arc.lock().unwrap(); if let Some(val) = rx.next().await { *data = val } })); } futures::future::select_all(handlers).await; } } fn main() { let mut svc = Service { handlers: Vec::new(), data: Arc::new(Mutex::new(true)), }; let (tx, rx) = mpsc::channel(8); svc.handlers.push(rx); let jh = task::spawn(async move { svc.handle_requests().await }); task::block_on(async move { tx.send(false).await; jh.await }); }
运行后出现如下错误:
error[E0277]: `std::sync::MutexGuard<'_, bool>` cannot be sent between threads safely --> src/main.rs:34:14 | 34 | let jh = task::spawn(async move { svc.handle_requests().await }); | ^^^^^^^^^^^ `std::sync::MutexGuard<'_, bool>` cannot be sent between threads safely
我理解handle_requests生成的Future因MutexGuard不满足Send trait,不符合spawn的约束,但不知如何实现该逻辑。我需要多个可修改data的Future,本以为这正是Mutex的适用场景。如何让Future修改共享数据且不违反spawn的trait要求?或者spawn并非合适选择?
编辑1
我将代码片段从
let mut data = data_arc.lock().unwrap(); if let Some(val) = rx.next().await { *data = val }
改为
if let Some(val) = rx.next().await { *data_arc.lock().unwrap() = val }
后错误消失,但不清楚原因。
编辑2
我认为这是典型的“跨await边界持有mutex guard”问题,只是被类型检查器混淆了。
问题解析与解决方案
问题根源
你遇到的确实是跨await边界持有MutexGuard的问题。当你在await之前获取MutexGuard(let mut data = data_arc.lock().unwrap();),这个Guard会被保存在Future的挂起状态中。而std::sync::MutexGuard并不实现Send trait——它绑定了当前线程的生命周期,跨线程传递可能导致未定义行为,但async_std::task::spawn要求传入的Future必须是Send的(任务可能被调度到不同线程执行),因此编译器抛出错误。
把获取Guard的操作移到await之后,MutexGuard的生命周期就被限制在await后的同步代码块里,不会被包含在Future的挂起状态中。此时整个Future不再持有跨线程不安全的MutexGuard,自然满足Send约束,错误也就消失了。
更优的异步共享数据方案
虽然修改后的代码能运行,但std::sync::Mutex是为同步场景设计的,在异步代码中使用容易阻塞线程。更合适的选择是用异步专用的Mutex:
- 如果使用Tokio runtime,用
tokio::sync::Mutex - 如果使用async_std,用
async_std::sync::Mutex
这些异步Mutex的Guard实现了Send,允许跨await边界持有,且不会阻塞整个线程。用async_std::sync::Mutex改造你的代码示例:
use async_std::task; use async_std::sync::Mutex; // 替换为异步Mutex use futures::channel::mpsc::{self, Receiver}; use futures::{SinkExt, StreamExt}; use std::sync::Arc; struct Service { handlers: Vec<Receiver<bool>>, data: Arc<Mutex<bool>>, } impl Service { async fn handle_requests(&mut self) { let mut handlers = Vec::new(); for rx in &mut self.handlers { let data_arc = self.data.clone(); handlers.push(Box::pin(async move { if let Some(val) = rx.next().await { let mut data = data_arc.lock().await; // 异步lock,Guard是Send的 *data = val; } })); } futures::future::select_all(handlers).await; } } fn main() { let mut svc = Service { handlers: Vec::new(), data: Arc::new(Mutex::new(true)), }; let (tx, rx) = mpsc::channel(8); svc.handlers.push(rx); let jh = task::spawn(async move { svc.handle_requests().await }); task::block_on(async move { tx.send(false).await; jh.await }); }
这样既保证了异步场景下的线程安全,又符合spawn的Send要求,同时避免了不必要的线程阻塞。
内容的提问来源于stack exchange,提问作者Ben Little

