Tokio mpsc将Sender设为static时通道意外关闭的原因
问题分析:Tokio MPSC通道发送失败提示已关闭的原因
问题场景
你编写的Rust代码逻辑如下:
- 创建Tokio无界MPSC通道,将Sender存入全局
OnceLock - 构建新的Tokio Runtime并spawn异步任务,任务持有Receiver并循环轮询消息
- 主线程通过全局Sender发送消息,但收到通道已关闭的
SendError,即使添加std::thread::sleep也无法解决
代码片段:
use tokio::sync::mpsc::{error::SendError, unbounded_channel, UnboundedSender}; use std::sync::OnceLock; static GSENDER: OnceLock<UnboundedSender<&'static str>> = OnceLock::new(); fn main() { let (sender, mut channel) = unbounded_channel(); GSENDER.set(sender).unwrap(); tokio::runtime::Builder::new_multi_thread() .worker_threads(1) // on a new thread .enable_all() .build() .unwrap() .spawn(async move { println!("[{:?}] Starting channel", chrono::Utc::now()); while let Some(msg) = channel.recv().await { println!("[{:?}] Recvd: {msg}", chrono::Utc::now()); } println!("[{:?}] Closing channel", chrono::Utc::now()); }); // Does not help, as it shouldn't anyway // std::thread::sleep(std::time::Duration::from_secs(1)); if let Some(channel_in) = GSENDER.get() { if let Err(SendError(_)) = channel_in.send("test") { println!("[{:?}] Channel down", chrono::Utc::now()); } } else { unreachable!() } }
核心原因
问题出在Tokio Runtime的生命周期管理上:
- 你通过
Builder::build()创建的Runtime是一个临时值,调用spawn后没有被任何变量持有,语句执行完毕后立即被销毁(drop) - Runtime被销毁时,会强制终止所有它管理的异步任务,此时持有Receiver的任务甚至还没来得及启动执行,Receiver就被随之销毁
- Tokio MPSC通道的规则是:当所有Receiver实例被销毁后,通道会进入关闭状态,此时Sender发送消息必然返回
SendError
添加std::thread::sleep无效的原因是:Runtime已经被销毁,里面的任务根本没有机会运行,Receiver早已被释放。
解决方案
需要确保Runtime的生命周期覆盖任务执行和消息发送的全流程,以下是两种可行写法:
写法1:持有Runtime并等待任务完成
use tokio::sync::mpsc::{error::SendError, unbounded_channel, UnboundedSender}; use std::sync::OnceLock; use tokio::task::JoinHandle; static GSENDER: OnceLock<UnboundedSender<&'static str>> = OnceLock::new(); fn main() { let (sender, mut channel) = unbounded_channel(); GSENDER.set(sender).unwrap(); // 持有Runtime实例,避免提前销毁 let runtime = tokio::runtime::Builder::new_multi_thread() .worker_threads(1) .enable_all() .build() .unwrap(); // 保存任务句柄,用于后续等待 let handle: JoinHandle<_> = runtime.spawn(async move { println!("[{:?}] Starting channel", chrono::Utc::now()); while let Some(msg) = channel.recv().await { println!("[{:?}] Recvd: {msg}", chrono::Utc::now()); } println!("[{:?}] Closing channel", chrono::Utc::now()); }); if let Some(channel_in) = GSENDER.get() { if let Err(SendError(_)) = channel_in.send("test") { println!("[{:?}] Channel down", chrono::Utc::now()); } } else { unreachable!() } // 等待任务执行完成,保证Runtime不会提前销毁 runtime.block_on(handle).unwrap(); }
写法2:使用block_on管理异步逻辑(更推荐)
直接用Runtime的block_on方法运行整个异步逻辑,避免手动管理生命周期:
use tokio::sync::mpsc::{error::SendError, unbounded_channel, UnboundedSender}; use std::sync::OnceLock; static GSENDER: OnceLock<UnboundedSender<&'static str>> = OnceLock::new(); fn main() { tokio::runtime::Builder::new_multi_thread() .worker_threads(1) .enable_all() .build() .unwrap() .block_on(async { let (sender, mut channel) = unbounded_channel(); GSENDER.set(sender).unwrap(); let handle = tokio::spawn(async move { println!("[{:?}] Starting channel", chrono::Utc::now()); while let Some(msg) = channel.recv().await { println!("[{:?}] Recvd: {msg}", chrono::Utc::now()); } println!("[{:?}] Closing channel", chrono::Utc::now()); }); if let Some(channel_in) = GSENDER.get() { if let Err(SendError(_)) = channel_in.send("test") { println!("[{:?}] Channel down", chrono::Utc::now()); } } else { unreachable!() } handle.await.unwrap(); }); }
内容的提问来源于stack exchange,提问作者porkbrain
相关产品推荐
相关产品推荐

