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:使用通道迭代器自动接收所有消息
如果不想硬编码任务数量,可以用通道的迭代器,它会持续阻塞接收消息,直到所有发送端都被销毁才会终止迭代,需要做两个调整:
- 提交完所有任务后,手动销毁入口通道的发送端,触发所有 worker 执行完当前任务后收到Halt消息退出
- 用迭代器遍历出口通道
#[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
相关产品推荐
相关产品推荐

