如何限制FuturesUnordered并行处理数且支持动态添加任务?
问题
我正在用FuturesUnordered把异步任务加入多线程tokio运行时的队列,这些异步任务会返回不同类型的结果,所以我把每个任务的结果映射成了自定义的Event枚举类型。现有代码运行正常,但我想限制FuturesUnordered同时并行处理的任务数量。尝试用buffered(10)时出现了编译错误(提示Event未实现Future trait),请问怎么实现限制并行数同时还能动态添加任务?
原代码示例
enum Event { ResultTypeA {...}, ResultTypeB {...}, ResultTypeC {...}, ResultTypeD {...} } let pending_futures: FuturesUnordered<Pin<Box<dyn Future<Output = Event> + Send>>> = FuturesUnordered::default() loop { tokio::select! { Some(future) = workload_receiver.recv() => { pending_futures.push(future.boxed()); }, Some(event) = pending_futures.next() => process_event(event), else => break, } }
报错的修改代码
loop { tokio::select! { Some(future) = workload_receiver.recv() => { pending_futures.push(future.boxed()); }, Some(event) = pending_futures.buffered(10).next() => process_event(event), else => break, } }
编译错误信息
--> src/main.rs | 257 | Some(event) = pending_futures.buffered(10).next() => process_event(event), | ^^^^^^^^^^^^ `Event` 未实现 `Future` trait | = help: 特征 `futures::Future` 未为 `Event` 实现 = note: Event 必须是 future 或实现 `IntoFuture` 才能被等待 note: 由 `buffered` 中的约束所要求 --> futures-util-0.3.24/src/stream/stream/mod.rs:1359:21 | 1359 | Self::Item: Future, | ^^^^^^ 此约束在 `buffered` 中被要求
解决方案
错误原因
buffered(10)的作用是处理元素为Future的Stream,它会并行执行最多10个该Stream中的Future,要求Stream的Item必须是Future类型。但你的FuturesUnordered本身就是一个输出Event的Stream(它内部已经执行完Future并返回结果了),所以调用buffered会触发类型错误,因为Event不是Future。
正确实现:用信号量控制并发
要限制并行任务数量,同时支持动态添加任务,最直接的方式是用Tokio的Semaphore来做并发控制。每个异步任务在执行前先获取信号量的许可,执行完毕后自动释放许可,这样就能保证同时运行的任务数不超过设定的上限。
修改后的代码示例:
use tokio::sync::Semaphore; use std::sync::Arc; enum Event { ResultTypeA, ResultTypeB, ResultTypeC, ResultTypeD, } // 创建信号量,设置最大并发数为10 let semaphore = Arc::new(Semaphore::new(10)); let pending_futures: FuturesUnordered<Pin<Box<dyn Future<Output = Event> + Send>>> = FuturesUnordered::default(); loop { tokio::select! { Some(future) = workload_receiver.recv() => { let semaphore = semaphore.clone(); // 把任务包装成先获取许可再执行的Future let wrapped_future = async move { // 获取许可,无可用许可时会阻塞等待 let _permit = semaphore.acquire().await.unwrap(); // 执行原任务并返回Event future.await }; pending_futures.push(wrapped_future.boxed()); }, Some(event) = pending_futures.next() => process_event(event), else => break, } }
原理说明
- Semaphore初始化时设置10个许可,代表最多同时有10个任务能获取到许可。
- 每个新任务被包装成先获取许可再执行的Future,当没有可用许可时,
semaphore.acquire().await会阻塞,直到有任务执行完毕释放许可。 - 许可
_permit在Future执行完毕后会自动被Drop,从而释放回信号量,无需手动处理。 - 这种方式既保留了FuturesUnordered动态添加任务的能力,又严格限制了并行执行的任务数量。
内容的提问来源于stack exchange,提问作者Basic Coder
相关产品推荐
相关产品推荐

