如何在Tokio TCP服务器达连接上限时停止接收新TCPStream
解决Tokio TCP服务器连接数限制且不阻塞其他异步任务的问题
核心问题分析
你原来的写法会阻塞data_printing任务,原因在于:
- 使用同步Mutex做连接数计数,循环中反复调用
lock()会占用Tokio调度器的时间片,导致其他异步任务无法被调度。 - 若连接数已满,循环会持续轮询判断条件,无法让出线程给其他任务,最终阻塞整个
join!中的所有任务。
正确解决方案:使用Tokio异步信号量(Semaphore)
Tokio提供的tokio::sync::Semaphore是专门用于异步场景的限流工具,当没有可用许可时,acquire_owned().await会挂起当前任务,主动让出线程,让调度器处理其他任务,完全不会阻塞data_printing这类并行任务。
修改步骤
- 替换计数工具:把
Arc<Mutex<u8>>替换为Arc<tokio::sync::Semaphore>,初始化时设置最大许可数为6(即最大连接数)。 - 获取许可再接受连接:在调用
listener.accept()前,先获取Semaphore的许可,确保只有在有可用连接名额时才接受新连接。 - 自动释放许可:获取的
Permit会在任务结束后自动drop释放许可,无需手动维护连接数计数。
修改后的代码示例
监听函数修改
use tokio::sync::Semaphore; use std::sync::Arc; use tokio::net::TcpListener; use tokio::io::BufReader; use std::sync::Mutex; async fn listen(connected_semaphore: Arc<Semaphore>, tried_numbers: Arc<Mutex<Vec<u32>>>) { let listener = TcpListener::bind("localhost:8881").await.unwrap(); loop { // 获取连接许可,无可用时挂起任务,让出线程 let permit = connected_semaphore.acquire_owned().await.unwrap(); let (socket, _) = listener.accept().await.unwrap(); let tried_numbers = tried_numbers.clone(); // 将许可传入任务,任务结束后自动释放 tokio::spawn(async move { let (sread, _swrite) = socket.split(); let reader = BufReader::new(sread); process_socket(tried_numbers, reader).await; // permit会在此处自动drop,释放连接名额 }); } }
main函数初始化调整
// 初始化信号量,最大连接数6 let connected_semaphore = Arc::new(Semaphore::new(6)); let tried_numbers = Arc::new(Mutex::new(Vec::new())); tokio::join!( listen(connected_semaphore.clone(), tried_numbers.clone()), data_printing(connected_semaphore.clone(), tried_numbers.clone(), total_unique_numbers, output_file) );
额外说明
- 若需获取当前连接数,可通过
6 - connected_semaphore.available_permits()计算。 - 异步代码中避免使用同步Mutex的
lock().unwrap(),这类同步阻塞操作会破坏Tokio的异步调度模型,优先使用Tokio提供的异步同步原语(如Semaphore、tokio::sync::Mutex)。
内容的提问来源于stack exchange,提问作者19mike95
相关产品推荐
相关产品推荐

