Tokio中实现Go ErrGroup等效并发模式的简洁方案探究
Tokio中实现Go ErrGroup等效并发模式的简洁方案探究
我从Go转到Tokio开发时,特别怀念ErrGroup这个工具——它完美封装了一个很常见的并发模式:并行执行多个任务,只要其中一个失败,就取消所有其余任务,阻塞直到所有任务都终止,最后返回第一个出现的错误。
先给大家看个Go里的示例代码:
func runTasks() error { g, ctx := errgroup.WithContext(context.Background()) for id := 1; id <= 3; id++ { g.Go(func() error { select { // 模拟任务工作 case <-time.After(time.Millisecond * 500 * time.Duration(id)): if id == 2 { return errors.New("simulated error") } return nil case <-ctx.Done(): return context.Cause(ctx) } }) } return g.Wait() }
那在Tokio(或者说异步Rust)里,有没有更简洁的实现方式呢?我目前想到的写法是这样的:
use anyhow::{bail, Result}; use tokio::task::JoinSet; use tokio::time::Duration; use tokio_util::sync::CancellationToken; #[tokio::main] async fn main() -> Result<()> { let mut join_set = JoinSet::new(); let cancel_token = CancellationToken::new(); for id in 1..=3 { let child_token = cancel_token.child_token(); join_set.spawn(async move { tokio::select! { // 模拟任务工作 _ = tokio::time::sleep(Duration::from_millis(500) * id) => { if id == 2 { bail!("simulated error") } else { Ok(()) } } _ = child_token.cancelled() => bail!("canceled"), } }); } let mut final_result: Result<()> = Ok(()); while let Some(result) = join_set.join_next().await { match result.expect("join error") { Ok(_) => println!("Task completed successfully"), Err(e) => { println!("Task failed with error: {:?}", e); if final_result.is_ok(){ final_result = Err(e) } cancel_token.cancel(); break; } } } // 等待剩余任务因取消而退出 let _ = join_set.join_all().await; final_result }
我感觉基于现有Tokio的功能,完全可以封装一个类似ErrGroup的小工具库,但我还是好奇:有没有现成的实现我没发现?或者我上面的写法有没有什么潜在问题?
备注:内容来源于stack exchange,提问作者Brendan
相关产品推荐
相关产品推荐

