Rust中tokio::spawn调用futures::stream::unfold的Send/Sync问题
解决Tokio Spawn中Stream的Send trait问题
问题背景
我仅做过Rust小型练手项目,当前需处理从Redis消费数据并与其他服务交互的逻辑。定义了Protocol trait并基于fred 5.2.0实现Redis传输的subscribe方法,直接运行代码正常,但使用tokio::spawn生成任务时触发Send/Sync相关错误。
定义的Protocol trait
#[async_trait] pub trait Protocol { async fn send(&self, value: &str) -> Result<()>; async fn subscribe(&self) -> Result<Pin<Box<dyn Stream<Item = String>>>>; }
实现的subscribe方法
async fn subscribe(&self) -> Result<Pin<Box<dyn Stream<Item = String>>>> { let c = self.client.clone(); let key = self.input_key.clone(); let s = stream::unfold((c, key), |(c, key)| async move { loop { let value = match c.blpop::<(String, String), _>(&key, 2.0).await { Ok((_, value)) => value, _ => { continue; } }; return Some((value, (c, key))); } }); Ok(Box::pin(s)) }
可行代码(直接运行)
pub async fn run(&self) { let protocol = self.protocol.clone(); let s = protocol.subscribe().await.unwrap(); s.for_each(|x| async move { println!("Received {}", x); }).await; }
报错代码(使用tokio::spawn)
pub async fn run(&self) { let protocol = self.protocol.clone(); tokio::spawn(async move { let s = protocol.subscribe().await.unwrap(); s.for_each(|x| async move { println!("Received {}", x); }).await }); }
错误信息
error: future cannot be sent between threads safely --> src/foo.rs:24:22 | 24 | tokio::spawn(async move { | ______________________^ 25 | | let s = protocol.subscribe().await.unwrap(); 26 | | s.for_each(|x| async move { 27 | | println!("Received {}", x); 28 | | }).await 29 | | }); | |_________^ future created by async block is not `Send` | = help: the trait `std::marker::Send` is not implemented for `dyn futures::Stream<Item = std::string::String>` note: future is not `Send` as this value is used across an await --> src/foo.rs:28:15 | 25 | let s = protocol.subscribe().await.unwrap(); | - has type `Pin<Box<dyn futures::Stream<Item = std::string::String>>>` which is not `Send` ... 28 | }).await | ^^^^^^ await occurs here, with `s` maybe used later 29 | }); | - `s` is later dropped here note: required by a bound in `tokio::spawn` --> /home/bam/.cargo/registry/src/github.com-1ecc6299db9ec823/tokio-1.26.0/src/task/spawn.rs:163:21 | 163 | T: Future + Send + 'static, | ^^^^ required by this bound in `tokio::spawn`
解决步骤
1. 为trait返回的Stream添加Send约束
修改Protocol trait的subscribe方法,明确要求返回的Stream实现Send trait:
#[async_trait] pub trait Protocol { async fn send(&self, value: &str) -> Result<()>; // 添加Send约束,确保Stream可跨线程传递 async fn subscribe(&self) -> Result<Pin<Box<dyn Stream<Item = String> + Send>>>; }
2. 确认实现的Stream满足Send要求
fred 5.x的Client本身实现了Send + Sync,input_key是String也满足Send,因此stream::unfold生成的天然是Send类型的Stream,修改trait约束后就能匹配要求。
3. 额外注意事项
- 若需要在多线程间传递
Protocoltrait对象,需确保trait对象满足Send + Sync,比如使用Box<dyn Protocol + Send + Sync>。 - 确保
tokio::spawn闭包中的所有变量满足'static约束,这里通过cloneprotocol已满足要求。
内容的提问来源于stack exchange,提问作者Benjamin Podszun
相关产品推荐
相关产品推荐

