如何竞速多个ZIO Effect并获取最先完成的n个结果?
实现ZIO的
raceN方法:获取前N个完成的结果并中断其余任务 要实现raceN方法,核心是并行运行所有任务,收集最先完成的n个结果,一旦达标立即中断剩余未完成任务。我们可以借助ZIO的Queue控制结果收集,结合Fiber中断机制实现需求,具体代码如下:
import zio._ final def raceN[R, E, A](n: Int)(as: Iterable[ZIO[R, E, A]]): ZIO[R, Nothing, List[Either[E, A]]] = { // 前置参数校验:n必须在1到任务总数之间 require(n > 0 && n <= as.size, "n必须大于0且不超过任务数量") for { // 创建容量为n的有界队列,仅保留前n个完成的结果 resultQueue <- Queue.bounded[Either[E, A]](n) // 包装每个任务:执行完成后将结果(捕获错误为Either)入队,忽略队列满后的入队失败 wrappedTasks = as.map(task => task.either.flatMap(resultQueue.offer).ignore) // 并行启动所有包装后的任务,生成可中断的Fiber runnerFiber <- ZIO.foreachPar(wrappedTasks)(identity).fork // 从队列中取出n个结果,直到收集完成 collectedResults <- ZIO.collectAll(List.fill(n)(resultQueue.take)) // 中断所有未完成的任务 _ <- runnerFiber.interrupt } yield collectedResults }
关键逻辑说明
- 队列控量:用
Queue.bounded(n)限制仅收集n个结果,后续任务的入队操作会因队列满失败,通过ignore忽略该失败,不影响任务正常终止。 - 并行执行与中断:通过
ZIO.foreachPar并行启动所有任务,并将其包装为Fiber;一旦收集到n个结果,立即中断该Fiber,终止所有未完成的任务。 - 错误捕获:用
task.either将任务的成功/失败结果统一转为Either[E,A],确保方法返回ZIO[R, Nothing, ...],不会因单个任务失败而终止整个流程。
应用场景示例:多数据源结果校验
假设我们从多个数据源并行加载同一数据,需要取最先返回的2个结果验证一致性:
// 模拟不同响应速度的数据源 val source1: ZIO[Any, Throwable, String] = ZIO.succeed("user_123_info").delay(120.millis) val source2: ZIO[Any, Throwable, String] = ZIO.succeed("user_123_info").delay(90.millis) val source3: ZIO[Any, Throwable, String] = ZIO.succeed("invalid_user_info").delay(60.millis) // 执行raceN并校验结果 val validateData = raceN(2)(List(source1, source2, source3)).flatMap { results => // 提取成功返回的结果 val validResults = results.collect { case Right(data) => data } validResults.size match { case 2 if validResults.distinct.size == 1 => ZIO.succeed(validResults.head) case 2 => ZIO.fail(new RuntimeException("多数据源返回结果不一致")) case _ => ZIO.fail(new RuntimeException("成功返回的结果不足2个")) } }
内容的提问来源于stack exchange,提问作者fast tooth
相关产品推荐
相关产品推荐

