Tokio Broadcast Channel Receiver无法接收消息问题求助
使用Tokio广播通道时无法收到消息的问题
我正在尝试用Tokio和Broadcast Channel编写服务器与客户端程序,逻辑是监听连接、读取TcpStream后通过通道发送消息。运行代码时,每次连接服务器并读取字节都会打印对应信息,但始终看不到Received的输出。
我的代码
use dbjade::serverops::ServerOp; use tokio::io::{BufReader}; use tokio::net::TcpStream; use tokio::{net::TcpListener, io::AsyncReadExt}; use tokio::sync::broadcast; const ADDR: &str = "localhost:7676"; // Your own address : TODO change to be configured const CHANNEL_NUM: usize = 100; use std::io; use std::net::{SocketAddr}; use bincode; #[tokio::main] async fn main() { // Create listener instance that bounds to certain address let listener = TcpListener::bind(ADDR).await.map_err(|err| panic!("Failed to bind: {err}")).unwrap(); let (tx, mut rx) = broadcast::channel::<(ServerOp, SocketAddr)>(CHANNEL_NUM); loop { if let Ok((mut socket, addr)) = listener.accept().await { let tx = tx.clone(); let mut rx = tx.subscribe(); println!("Receieved stream from: {}", addr); let mut buf = vec![0, 255]; tokio::select! { result = socket.read(&mut buf) => { match result { Ok(res) => println!("Bytes Read: {res}"), Err(_) => println!(""), } tx.send((ServerOp::Dummy, addr)).unwrap(); } result = rx.recv() =>{ let (msg, addr) = result.unwrap(); println!("Receieved: {msg}"); } } } } }
问题原因
- 单次select执行后退出:当前代码在接受连接后仅执行一次
tokio::select!,要么处理一次读取操作,要么处理一次广播接收,执行完成后就退出该连接的处理逻辑,无法持续监听后续的读取或广播消息。 - 未异步 spawn 连接任务:主循环处理连接时没有将逻辑放入独立异步任务,导致主循环被阻塞,无法同时处理多个连接,且每个连接只能处理一次操作。
- 广播订阅者生命周期过短:每个连接创建的订阅者
rx仅在单次select中存在,执行完就被丢弃,无法接收后续的广播消息。 - 缓冲区初始化错误:
vec![0, 255]实际创建了仅包含两个元素的数组,无法正常读取大量数据。
修复方案
将每个连接的处理逻辑放入独立异步任务,并在任务内部用循环包裹select!,持续监听读取和广播事件:
use dbjade::serverops::ServerOp; use tokio::io::AsyncReadExt; use tokio::net::{TcpListener, TcpStream}; use tokio::sync::broadcast; const ADDR: &str = "localhost:7676"; const CHANNEL_NUM: usize = 100; use std::net::SocketAddr; #[tokio::main] async fn main() { let listener = TcpListener::bind(ADDR) .await .unwrap_or_else(|err| panic!("Failed to bind: {err}")); let (tx, _) = broadcast::channel::<(ServerOp, SocketAddr)>(CHANNEL_NUM); loop { if let Ok((socket, addr)) = listener.accept().await { let tx = tx.clone(); // 为每个连接启动独立异步任务 tokio::spawn(async move { println!("Received stream from: {}", addr); let mut socket = socket; let mut rx = tx.subscribe(); let mut buf = vec![0; 256]; // 修正为256字节的缓冲区 loop { tokio::select! { result = socket.read(&mut buf) => { match result { Ok(0) => { // 连接正常关闭 println!("Connection closed: {}", addr); break; } Ok(res) => { println!("Bytes Read: {res}"); // 发送广播消息并处理错误 if let Err(e) = tx.send((ServerOp::Dummy, addr)) { eprintln!("Failed to send broadcast: {e}"); } } Err(e) => { eprintln!("Read error from {}: {}", addr, e); break; } } } result = rx.recv() => { match result { Ok((msg, sender_addr)) => { println!("Received: {msg} from {}", sender_addr); } Err(broadcast::RecvError::Closed) => { eprintln!("Broadcast channel closed"); break; } Err(broadcast::RecvError::Lagged(n)) => { eprintln!("Lagged behind by {} messages", n); } } } } } }); } } }
关键修改点
- 用
tokio::spawn将每个连接的处理逻辑独立为异步任务,主循环可继续接受新连接,实现多连接并发处理。 - 用
loop包裹select!,持续监听读取和广播事件,直到连接关闭或通道出错。 - 修正缓冲区初始化错误,改为创建256字节的空缓冲区。
- 补充完整错误处理,覆盖连接关闭、读取失败、广播通道异常等场景。
内容的提问来源于stack exchange,提问作者DrProfSgtMrJ
相关产品推荐
相关产品推荐

