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

Rust扩展The Book线程池时egress mpsc通道仅收到首条消息问题

问题根因

测试代码中的接收循环逻辑错误:
你写的 for s in tpool.egress_rx.recv() 中,recv() 是单次接收方法,返回值类型为 Result<String, RecvError>,该类型实现了 IntoIterator 特性,仅会产出1个元素,因此遍历完第一个返回的字符串就直接退出循环,后续消息完全没有处理。

修复方案

有两种常用的修复方式:

方案1:固定次数接收

你明确提交了10个任务,直接循环10次调用接收方法即可:

#[test]
fn tp_test() {
    let tpool = ThreadPool::new(4).expect("Cannot create threadpool");
    
    for i in 0..10 {
        let closure = move |worker_id| run_me(worker_id, i);
        tpool.execute(closure);
    }
    for _ in 0..10 {
        let s = tpool.egress_rx.recv().unwrap();
        println!("{}", s);
    }        
}

方案2:使用通道迭代器自动接收所有消息

如果不想硬编码任务数量,可以用通道的迭代器,它会持续阻塞接收消息,直到所有发送端都被销毁才会终止迭代,需要做两个调整:

  1. 提交完所有任务后,手动销毁入口通道的发送端,触发所有 worker 执行完当前任务后收到Halt消息退出
  2. 用迭代器遍历出口通道
#[test]
fn tp_test() {
    let tpool = ThreadPool::new(4).expect("Cannot create threadpool");
    
    for i in 0..10 {
        let closure = move |worker_id| run_me(worker_id, i);
        tpool.execute(closure);
    }
    // 销毁入口发送端,触发worker退出流程
    drop(tpool.ingress_tx);
    // 迭代器持续收消息直到所有出口发送端随worker销毁自动终止
    for s in tpool.egress_rx.iter() {
        println!("{}", s);
    }        
}

注:如果ingress_tx是ThreadPool的私有字段,你可以给ThreadPool新增一个公开的关闭方法,内部执行drop逻辑即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 14:09:03