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

将引用传递给生成Tokio任务的函数时遇生命周期错误求解

Rust提取异步函数时遇到'static生命周期报错的解决方法

问题背景

原本运行正常的WebSocket批量连接监听代码,在将核心逻辑提取为独立异步函数后,编译器抛出了关于'static生命周期的错误。

原正常运行代码

pub async fn start(num_websockets: i32) {
    let messages_received = Arc::new(Mutex::new(0));
    let (shutdown_signal, _) = broadcast::channel(1);
    let (done, mut wait_for_tasks) = mpsc::channel::<&str>(1);
    tokio::spawn(stats::message_counter(
        Arc::clone(&messages_received),
        shutdown_signal.subscribe(),
        done.clone(),
    ));
    for _i in 0..num_websockets {
        let (ws_stream, _) = websocket::connect().await.unwrap();
        let (mut write, read) = ws_stream.split();
        subscription::subscribe(&mut write, "book.BTC-PERPETUAL.100ms")
            .await
            .unwrap();
        tokio::spawn(websocket::listener(
            read,
            messages_received.clone(),
            shutdown_signal.subscribe(),
            done.clone(),
        ));
    }
}

提取函数后的代码

// start函数内的循环部分
for _i in 0..num_websockets {
    subscribe_and_listen(messages_received.clone(), shutdown_signal.subscribe(), done.clone()).await;
}

// 提取出的异步函数
async fn subscribe_and_listen(
    messages_received: Arc<Mutex<i32>>,
    shutdown_signal: broadcast::Receiver<&str>,
    done: mpsc::Sender<&str>,
) {
    let (ws_stream, _) = websocket::connect().await.unwrap();
    let (mut write, read) = ws_stream.split();
    subscription::subscribe(&mut write, "book.BTC-PERPETUAL.100ms")
        .await
        .unwrap();
    tokio::spawn(
        websocket::listener(
            read,
            messages_received.clone(),
            shutdown_signal,
            done.clone(),
        )
    );
}

编译器错误信息

error[E0759]: `shutdown_signal` has an anonymous lifetime `'_` but it needs to satisfy a `'static` lifetime requirement
   --> src/benchmark.rs:21:9
    |
13  |                                 shutdown_signal: broadcast::Receiver<&str>,
    |                                                  ------------------------- this data with an anonymous lifetime `'_`...
...
21  | /         websocket::listener(
22  | |             read,
23  | |             messages_received.clone(),
24  | |             shutdown_signal,
    | |             --------------- ...is used here...
25  | |             done.clone(),
26  | |         )
    | |_________^
    |
note: ...and is required to live as long as `'static` here
   --> src/benchmark.rs:20:5
    |
20  |     tokio::spawn(
    |     ^^^^^^^^^^^^
note: `'static` lifetime requirement introduced by this bound
   --> /Users/b0nes/.cargo/registry/src/github.com-1ecc6299db9ec823/tokio-1.20.0/src/task/spawn.rs:127:28
    |
127 |         T: Future + Send + 'static,

问题原因

Tokio的tokio::spawn要求传入的Future必须满足**'static生命周期**——因为任务调度器无法保证任务的生命周期与调用spawn的函数一致,任务不能持有任何非'static的外部引用。

提取函数后,broadcast::Receiver<&str>和mpsc::Sender<&str>中的&str是引用类型,编译器推断其生命周期为匿名的'_,无法满足'static要求。原代码中编译器通过上下文推断出了刚好匹配的生命周期,但提取函数后这种上下文丢失,直接触发报错。

解决方法

将通道的消息类型从引用类型&str改为拥有所有权的String,让相关的Receiver和Sender不再依赖外部引用,自然满足'static生命周期要求。

修改后的代码

调整start函数的通道创建

pub async fn start(num_websockets: i32) {
    let messages_received = Arc::new(Mutex::new(0));
    // 使用String作为广播通道的消息类型
    let (shutdown_signal, _) = broadcast::channel::<String>(1);
    // 使用String作为mpsc通道的消息类型
    let (done, mut wait_for_tasks) = mpsc::channel::<String>(1);
    tokio::spawn(stats::message_counter(
        Arc::clone(&messages_received),
        shutdown_signal.subscribe(),
        done.clone(),
    ));
    for _i in 0..num_websockets {
        subscribe_and_listen(messages_received.clone(), shutdown_signal.subscribe(), done.clone()).await;
    }
}

调整subscribe_and_listen函数的参数类型

async fn subscribe_and_listen(
    messages_received: Arc<Mutex<i32>>,
    shutdown_signal: broadcast::Receiver<String>,
    done: mpsc::Sender<String>,
) {
    let (ws_stream, _) = websocket::connect().await.unwrap();
    let (mut write, read) = ws_stream.split();
    subscription::subscribe(&mut write, "book.BTC-PERPETUAL.100ms")
        .await
        .unwrap();
    tokio::spawn(
        websocket::listener(
            read,
            messages_received.clone(),
            shutdown_signal,
            done.clone(),
        )
    );
}

同步调整其他依赖函数

确保stats::message_counter和websocket::listener的参数类型也同步更新为broadcast::Receiver<String>和mpsc::Sender<String>,保持类型一致性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 10:54:20