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

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. 额外注意事项

  • 若需要在多线程间传递Protocol trait对象,需确保trait对象满足Send + Sync,比如使用Box<dyn Protocol + Send + Sync>。
  • 确保tokio::spawn闭包中的所有变量满足'static约束,这里通过cloneprotocol已满足要求。

内容的提问来源于stack exchange,提问作者Benjamin Podszun

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 22:07:44