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

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}");
                }
            }
        }
    }
}

问题原因

  1. 单次select执行后退出:当前代码在接受连接后仅执行一次tokio::select!,要么处理一次读取操作,要么处理一次广播接收,执行完成后就退出该连接的处理逻辑,无法持续监听后续的读取或广播消息。
  2. 未异步 spawn 连接任务:主循环处理连接时没有将逻辑放入独立异步任务,导致主循环被阻塞,无法同时处理多个连接,且每个连接只能处理一次操作。
  3. 广播订阅者生命周期过短:每个连接创建的订阅者rx仅在单次select中存在,执行完就被丢弃,无法接收后续的广播消息。
  4. 缓冲区初始化错误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 08:55:37