ZIO中两段看似等价的Future互操作代码为何运行结果不同?
错误信息解读
[info] java.util.concurrent.ExecutionException: Boxed Exception
[info] at scala.concurrent.impl.Promise$.scala$concurrent$impl$Promise$$resolve(Promise.scala:99)
[info] at scala.concurrent.impl.Promise$DefaultPromise.tryComplete(Promise.scala:278)
[info] at scala.concurrent.Promise.complete(Promise.scala:57)
[info] at scala.concurrent.Promise.complete$(Promise.scala:56)
[info] at scala.concurrent.impl.Promise$DefaultPromise.complete(Promise.scala:104)
[info] at scala.concurrent.Promise.failure(Promise.scala:109)
[info] at scala.concurrent.Promise.failure$(Promise.scala:109)
[info] at scala.concurrent.impl.Promise$DefaultPromise.failure(Promise.scala:104)
[info] at zio.Fiber.$anonfun$toFutureWith$2(Fiber.scala:334)
[info] at zio.internal.FiberContext.evaluateNow(FiberContext.scala:404)
[info] at zio.internal.FiberContext.$anonfun$evaluateLater$1(FiberContext.scala:787)
[info] at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
[info] at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
[info] at java.base/java.lang.Thread.run(Thread.java:829)
[info] Cause: java.lang.InterruptedException: Interrupted by fibers: #50
[info] at zio.Cause.$anonfun$squashWith$1(Cause.scala:403)
[info] at scala.Option.orElse(Option.scala:477)
[info] at zio.Cause.squashWith(Cause.scala:400)
[info] at zio.Cause.squashTraceWith(Cause.scala:425)
[info] at zio.Fiber.$anonfun$toFutureWith$2(Fiber.scala:334)
[info] at zio.internal.FiberContext.evaluateNow(FiberContext.scala:404)
[info] at zio.internal.FiberContext.$anonfun$evaluateLater$1(FiberContext.scala:787)
[info] at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
[info] at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
[info] at java.base/java.lang.Thread.run(Thread.java:829)
该错误核心是ZIO Fiber被中断,导致转换后的Future以失败告终,中断根源在于两段代码的Runtime使用逻辑和任务调度流程完全不同。
两段代码的核心差异
第一段代码:Runtime.default.unsafeRun(x.toFuture)
val x = ZStream.fromIterable(Iterable.empty).runDrain await(zio.Runtime.default.unsafeRun(x.toFuture))
执行流程与问题点:
x.toFuture仅完成ZIO任务到Scala Future的格式转换,不会触发任务执行,任务逻辑需要ZIO Runtime调度后才会运行。Runtime.default是临时Runtime实例,unsafeRun方法会在当前线程同步执行传入的ZIO逻辑,执行完成后会立即关闭自身线程池。- 后续Future尝试调度流处理任务时,Runtime线程池已被销毁,对应Fiber被迫中断,最终抛出
InterruptedException。
第二段代码:Runtime.global.unsafeRunToFuture(x)
val x = ZStream.fromIterable(Iterable.empty).runDrain await(zio.Runtime.global.unsafeRunToFuture(x))
执行流程与优势:
unsafeRunToFuture是ZIO官方提供的互操作专用API,兼具“触发任务调度”和“转换为Future”的双重作用,直接将ZIO任务提交到Runtime线程池运行。Runtime.global是全局共享的持久化Runtime,其线程池长期存活,不会在单次任务结束后关闭,因此任务可正常执行完成,Future顺利返回成功结果。
关键区别总结
toFuture是纯转换API,不触发执行;unsafeRunToFuture是执行+转换的组合API,是ZIO到Future互操作的标准方式。Runtime.default是临时Runtime,unsafeRun执行后自动销毁线程池;Runtime.global是全局常驻Runtime,线程池持续运行,适合跨任务复用。
内容的提问来源于stack exchange,提问作者Sign

