ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Finagle Futures 并发编程完全指南:从 FuturePool 到 flatMap / collect / join / select 的组合实战

Finagle Futures 并发编程完全指南:从 FuturePool 到 flatMap / collect / join / select 的组合实战 后端RPC框架【免费下载链接】finagleA fault tolerant, protocol-agnostic RPC system项目地址https://gitcode.com/gh_mirrors/fi/finagle点击查看免费下载Finagle 用com.twitter.util.Future作为贯穿客户端、服务端与过滤器链路的统一并发抽象把网络 RPC、磁盘读取、长计算等操作封装为可组合的异步值。本文以官方用户指南《Concurrent Programming with Futures》为主体结合 Finagle 源码 与 线程模型文档系统讲解为什么 Future 是轻量级线程、如何用 FuturePool 安全承载阻塞操作、如何用flatMap/collect/join/select完成顺序、并发与并行组合以及如何避免组合过程中引入死锁。读完你将掌握 Finagle 异步编程的核心范式并能直接写出可运行、可复制的 Future 组合代码。一、Future轻量级线程Finagle 使用 Future 来封装和组合并发操作如网络 RPC。Future 与线程有着直接的类比关系——它们提供了独立且重叠的控制流可以被视为featherweight threads轻量级线程。与操作系统线程不同Future 的构造成本极低传统线程那种必须谨慎控制数量的经济学在这里不成立当并发操作由 Future 表示时同时维持数百万个未完成的操作完全不成问题。Future 还将 Finagle 从操作系统与运行时线程调度器中解耦出来。这一点被 Finagle 用在重要场合例如利用线程偏置thread biasing来降低上下文切换开销。这个类比有一个必须牢记的告诫不要在 Future 中执行阻塞操作。Future 不是抢占式的它必须通过flatMap主动让出控制权。阻塞操作会破坏这种协作式调度阻止其他异步操作推进导致应用出现莫名的变慢、吞吐下降甚至死锁。当然阻塞操作与 Future 也有安全的结合方式下文会详细展开。从源码结构看Finagle 核心模块finagle-core中所有Service、过滤器与调度器均以com.twitter.util.Future为返回类型官方入口文档 index.rst 也明确写道Finagle 采用基于 Futures 的干净、简单、安全的并发编程模型。二、为什么阻塞操作是并发模型的天敌Finagle 与其他非阻塞事件驱动框架一样在同一个 JVM 进程内为所有客户端和服务器共享一个固定大小的 I/O worker 线程池ThreadingModel.rst。默认池大小相当保守每个逻辑 CPU 核 2 个线程下限为 8 个 worker。阻塞哪怕一个 I/O 线程就可能影响多个客户端和服务器。更隐蔽的一点是所有由收到消息触发的用户代码都运行在 Finagle I/O 线程上除非显式地在别处执行例如应用自己的线程池或 FuturePool。这包括服务器端的Service.apply和所有 Future 回调。由于 Twitter Futures 纯粹是协调机制、本身不描述任何执行环境回调由满足 Promise 的那个线程运行而通常正是 I/O 线程满足 Promise例如收到 RPC 响应所以回调也顺理成章地落在 I/O 线程上。这种跟随调用线程的设计减少了上下文切换代价是 I/O 线程对用户的阻塞代码或慢代码变得脆弱。不仅是阻塞 I/OJDBC 驱动、JDK File API会阻塞 Finagle 线程CPU 密集型计算同样危险——I/O 线程忙于处理应用级计算就意味着它们在怠慢关键的 RPC 事件。文档中的示例是 Scala 的permutations最坏情况 O(n!)无论把它放在服务器端Service.apply里还是放在客户端 Future 的map回调里都会饱和 I/O 线程import com.twitter.finagle.Service import com.twitter.finagle.http.{Request, Response} def process(client: Service[Request, Response]): Future[String] client(Request()).map(rep rep.contentString.permutations.mkString(\n))2.1 如何识别阻塞ThreadingModel.rst 给出了两个无需外部工具即可观测的指标blocking_ms累计在Await.result/Await.ready中阻塞 I/O 线程的总时间计数器。当它不为零时说明请求路径上存在阻塞。pending_io_events所有事件循环中排队待处理的 I/O 事件数 gauge。该指标攀升说明 I/O 队列堵塞、I/O 线程过载。请求路径上一般建议彻底去掉Await至于pending_io_events的合理值文档明确表示没有统一标准追求零或接近零可作为一种友好建议视工作负载不同两位数值也可能可接受。三、用 FuturePool 承载阻塞工作当你确实有阻塞性工作要做——例如同步风格的 I/O或某个不是用异步风格编写的库——应当使用com.twitter.util.FuturePool。FuturePool 管理一组不干别的活的专用线程因此阻塞操作不会拖停其他异步工作。源码佐证Finagle 的 DNS 解析器 DnsResolver 正是用FuturePool.unboundedPool承载阻塞式域名解析并在解析完成后用Future.collectToTry汇总结果。下面someIO是一个等待 I/O 并返回字符串的操作例如读文件。把someIO(): String包进FuturePool.unboundedPool会返回Future[String]从而可以安全地将这个阻塞操作与其他 Future 组合。Scalaimport com.twitter.util.{Future, FuturePool} def someIO(): String // does some blocking I/O and returns a string val futureResult: Future[String] FuturePool.unboundedPool { someIO() }Javaimport com.twitter.util.Future; import com.twitter.util.FuturePools; import static com.twitter.util.Function.func0; FutureString futureResult FuturePools.unboundedPool().apply( func0(() - someIO()); );3.1 全局卸载OffloadingFuturePool 也是 Finagle 把用户代码整体移出 I/O 线程的底层设施。从源码看OffloadFilter 在 Future 链中引入异步边界将 continuation 从 I/O 线程移入给定 FuturePool它的服务端实现还会用Promise.become正确传播中断避免直接打断 FuturePool 线程本身。其注释还记录了实测经验仅用service(request).flatMap(pool.apply(_))时约 6% 的情况会因竞态而卸载失败因此实现改为先创建Promise.interruptsRep再在池内更新把竞态失败率压到约 0.0001%。按 ThreadingModel.rst 的用法卸载可以按端点、按客户端/服务器、按整个 JVM 三种粒度进行import com.twitter.util.{Future, FuturePool} // 按端点方法粒度 def offloadedPermutations(s: String, pool: FuturePool): Future[String] pool(s.permutations.mkString(\n)) // 按整个客户端或服务器 import com.twitter.util.FuturePool import com.twitter.finagle.Http val server: Http.Server Http.server .withExecutionOffloaded(FuturePool.unboundedPool) val client: Http.Client Http.client .withExecutionOffloaded(FuturePool.unboundedPool)按整个应用JVM 进程启用全局卸载使用命令行标志-com.twitter.finagle.offload.autotrue或者手动调优线程配置-com.twitter.finagle.offload.numWorkers14 -com.twitter.finagle.netty4.numWorkers10对应的参数实现在 OffloadFuturePool.scala 与 numWorkers.scala 中numWorkers未显式指定且auto开启时会按com.twitter.jvm.numProcs().ceil推导池的默认兜底是FuturePool.unboundedPool。卸载类过滤器的行为在 OffloadFilterTest.scala 中有大量测试覆盖。开启全局卸载后还可选启用实验性的 Offload 准入控制根据工作队列等待时间拒绝工作默认阈值 20ms-com.twitter.finagle.offload.admissionControlenabled -com.twitter.finagle.offload.admissionControl50.milliseconds四、同步Synchronized与同步化Synchronous的区别同步synchronization与同步行为synchronous behavior是两回事。同步调用synchronous calls在同一个线程内等待某个工作完成后再执行下一条语句而同步化代码段synchronized sections则允许一个线程执行语句并在该线程完成包围的语句之前阻塞所有其他调用者。一个不加同步就会出错的经典例子def incrementAndReturn(): Integer { counter 1; counter }如果两个线程在同一个对象上并发执行incrementAndReturn有可能两个线程都在任一线程执行 return 之前执行了counter 1。一旦如此两个线程会得到相同的值例如 counter 初始为 45则两个线程都会拿到 47。保证每个线程都拿到唯一且不跳号的 counter 值最简单的办法是把临界语句包进 synchronized 块def incrementAndReturn(): Integer { this.synchronized { counter 1; counter } }4.1 锁对象的作用域synchronized 块的语法要求用户定义一个锁对象由它授予对后续代码块的访问权同一时刻只有一个线程能持有某个锁对象。上例以this为锁意味着类内任何this.synchronized {...}块都绑定到同一个对象this。例如def incrementAndReturn(): Integer { this.synchronized { counter 1; counter } } def decrementAndReturn(): Integer { this.synchronized { counter - 1; counter } }两个函数现在都被this门控不仅多个线程对incrementAndReturn会串行执行调用decrementAndReturn时也会在this上排队等待。类中其他方法可以通过省略 synchronized 块而自由执行、无需等待def incrementAndReturn(): Integer { this.synchronized { counter 1; counter } } def decrementAndReturn(): Integer { this.synchronized { counter - 1; counter } } def readCounter(): Integer { counter }这里任何线程都可以随时调用readCounter无需等待控制this。4.2 用专用锁对象细化同步粒度为了可读性和逻辑分段可以定义专门用作同步锁的对象而不是把整个实例作为锁的粒度private[this] var counter: Integer 0 private[this] val lock: Object counter def incrementAndReturn(): Integer { lock.synchronized { counter 1; counter } } def decrementAndReturn(): Integer { lock.synchronized { counter - 1; counter } } def readCounter(): Integer { counter }这给了我们演进类的灵活性随着类演化若发现新的需要同步的操作可以纳入lock对象的保护伞之下。当前lock与 counter 本身同义但将来可能换用其他成员作为锁或为不同的状态集合准备不同的锁。五、同步的风险活锁与死锁同步是定义临界区、让运行时管理阻塞/调度/控制权交接的有效语言特性能确保内部状态如上面的 counter可预测地变更。但它也把开发者暴露给一类新 bug线程无限期地等待数据变化或等待获取锁对象。两类常见的锁问题是活锁livelock与死锁deadlock。活锁线程都还活着但代码在等待某个数据变化才能继续。系统唤醒一个线程线程检查数据状态是否正确发现没有变化又睡回去。如果负责更新的进程/线程无法完成更新系统就陷入活锁。死锁两个或多个线程在某个 synchronized 语句上互相阻塞各自持有的锁对象正是对方等待的。例如一个流程中线程先独占访问某个 Person 对象再查询其兄弟siblings以便一起更新。若两个线程分别对一对兄弟执行该流程就可能死锁线程 A 持有 Person A 的锁线程 B 持有 Person A 的兄弟 Person B 的锁随后 A 等待 B 用完 Person BB 等待 A 用完 Person A——死锁。一个详细的真实死锁示例见下文组合中的同步与死锁。六、Future 作为容器三种状态与回调被 Future 表示的常见操作包括对远端主机的 RPC另一个线程中的长计算从磁盘读取注意这些操作都可能失败远端主机可能崩溃、计算可能抛异常、磁盘可能损坏。因此一个Future[T]恰好占据三种状态之一Emptypending尚未完成Succeeded以类型T的结果成功Failed携带一个Throwable虽然可以直接查询这个状态但很少有用。更常见的做法是注册回调在结果可用时接收它import com.twitter.util.Future val f: Future[Int] ??? f.onSuccess { res: Int println(The result is res) }上述回调只在成功时被调用。也可以注册处理失败的回调import com.twitter.util.Future val f: Future[Int] ??? f.onFailure { cause: Throwable println(f failed with cause) }七、顺序组合flatMap注册回调很有用但 API 比较笨拙。Future 的力量在于组合compose。大多数操作可以拆成更小的操作这些操作又构成复合操作Future 让创建这种复合操作变得容易。考虑一个典型的抓取网站缩略图的例子类似 Pinterest 的流程它通常包含三步抓取主页解析页面找到第一个图片链接抓取该图片链接这是顺序组合的典型场景要做下一步必须先成功完成上一步。在 Future 中这就是flatMap。flatMap的结果是一个代表该复合操作结果的 Future。假设有辅助方法fetchUrl抓取给定 URL和findImageUrls解析 HTML 页面找出图片链接实现如下import com.twitter.util.Future def fetchUrl(url: String): Future[Array[Byte]] ??? def findImageUrls(bytes: Array[Byte]): Seq[String] ??? val url https://www.google.com val f: Future[Array[Byte]] fetchUrl(url).flatMap { bytes val images findImageUrls(bytes) if (images.isEmpty) Future.exception(new Exception(no image)) else fetchUrl(images.head) } f.onSuccess { image println(Found image of size image.size) }f代表复合操作先取回网页再取回第一个图片链接。如果任一子操作失败第一次或第二次fetchUrl失败或者findImageUrls没找到任何图片复合操作也失败。细心的读者可能注意到这不就是分号semicolon的职责吗确实不远分号顺序执行两条语句对传统 I/O 操作而言效果与上面的flatMap相同异常机制扮演失败 Future 的角色。但 Future 灵活得多——组合可以进行并发与并行编排且失败传播、中断传播都可控这是分号无法做到的。术语注记flatMap这个名字的语源来自集合上的顺序组合与Future 上的顺序组合之间的深层类比即函数式编程中的 Monad 结构并不是随意命名。八、并发组合Future.collectFuture 也可以并发地组合。扩展上面的例子一次性抓取所有图片。并发组合由Future.collect提供import com.twitter.util.Future val collected: Future[Seq[Array[Byte]]] fetchUrl(url).flatMap { bytes val fetches findImageUrls(bytes).map { url fetchUrl(url) } Future.collect(fetches) }这里同时结合了并发与顺序组合先抓取网页然后并发收集所有底层图片的抓取结果。与顺序组合一样并发组合也会传播失败collected会在任一底层 Future 失败时失败。Future.collectToTry如果希望并发收集一组 Future 的同时累积错误而不是快速失败可以使用Future.collectToTry。Finagle 内部正是这么做的DNS 解析器 InetResolver.scala 对多个 host 的解析结果执行Future.collectToTry再统一合并成最终的Addr全部成功则Addr.Bound、全部未知主机则Addr.Neg、出现意外错误则Addr.Failed从而在部分成功时仍能利用成功的解析结果。另外编写自己的 Future 组合子也非常简单。这在分布式系统中带来了极大的模块化收益常见模式可以被干净地抽象出来。九、并行组合Future.join 的四种模式collect专门用于多次执行同一种操作、返回同类型结果、并想知道它们何时全部完成的场景。但我们常常会并行触发不同的操作返回不同类型的结果。当需要同时使用几种不同计算的输出时最好等所有结果都回来再继续。在经典线程编程模型里对应的做法是对 fork 出的线程调用join——这正是Future.join的用武之地。Future.join共有四种模式。它们最初为 Scala 编写但Future对象上的方法也有 Java 友好的版本Futures.joinFuture实例上的方法在 Java 中也可以直接使用。模式一Future#join实例方法。它接受另一个 Future 作为参数返回一个在this与参数都满足后才满足的 Future且包含两个 Future 的内容import com.twitter.util.Future val numFollowers: Future[Int] ??? val profileImageURL: Future[String] ??? val userProfileData: Future[(Int, String)] numFollowers.join(profileImageURL)模式二Future.join对象方法多个 Future。用于同时 join 多个不同结果有大量重载以支持不同数量的 Futureimport com.twitter.util.Future val numFollowers: Future[Int] ??? val profileImageURL: Future[String] ??? val followersYouKnow: Future[Seq[User]] ??? val userProfileData: Future[(Int, String, Seq[User])] Future.join(numFollowers, profileImageURL, followersYouKnow)模式三Future#joinWith。调用Future#join后常见的动作是立刻变换结果。作为小优化joinWith可以避免分配 Tuple2 实例import com.twitter.util.Future val numFollowers: Future[Int] ??? val profileImageURL: Future[String] ??? val constructUserProfile: (Int, String) UserProfile val userProfile: Future[UserProfile] numFollowers.joinWith(profileImageURL)(constructUserProfile)模式四Future.join(Seq[Future[A]]): Future[Unit]。这个比较特殊——像Future.collect一样作用于Seq[Future[A]]但最终只返回Future[Unit]即只告诉你所有组成部分是否全部成功。它被用来实现其他Future.join方法并作为一种小优化暴露给只需要知道成功或失败、不关心实际结果的使用场景import com.twitter.util.Future val numFollowers: Future[Int] ??? val profileImageURL: Future[String] ??? val followersYouKnow: Future[Seq[User]] ??? val profileDataIsReady: Future[Unit] Future.join(Seq(numFollowers, profileImageUrl, followersYouKnow))十、组合中的同步与死锁AsyncSemaphore 实例剖析如前所述同步会在并发环境中引入新一类 bug。文档给出一个真实世界的死锁实例com.twitter.util中AsyncSemaphore的修复提交。修复前fail(..)、release()与中断处理器interrupt handler在完成 Promise 时都同步于this这可能导致两个线程与两个独立的 AsyncSemaphore 交互时发生死锁。下面是一个隔离出该误行为的玩具示例——它看起来过于明显有问题而不会真实发生但恰好能展示可能意外出现的错误模式val semaphore1, semaphore2 new AsyncSemaphore(1) // The semaphores have already been taken: val permitForSemaphore1 await(semaphore1.acquire()) val permitForSemaphore2 await(semaphore2.acquire()) // The semaphores have had continuations attached as follows: semaphore1.acquire().flatMap { permit val otherWaiters semaphore2.numWaiters // synchronizing method permit.release() otherWaiters } semaphore2.acquire().flatMap { permit val otherWaiters semaphore1.numWaiters // synchronizing method permit.release() otherWaiters }现在触发死锁val threadOne new Thread { override def run() { permitForSemaphore1.release() } } val threadTwo new Thread { override def run() { permitForSemaphore2.release() } } threadOne.start threadTwo.start在这种情形下threadOne 与 threadTwo 可能死锁——仅从release()调用看并不明显。原因在于acquire()返回一个 Promise上面挂载了 continuation。当线程调用Permit.release()时AsyncSemaphore 实现会在锁对象this上同步并在退出 synchronized 块之前把 permit 交给下一个等待的 Promise。这就会执行 continuation而 continuation 调用了另一个 AsyncSemaphore 上的方法并试图再次同步——正如上文兄弟更新示例所描述的那样。实例semaphore1与semaphore2激进地锁住各自的 AsyncSemaphore然后阻塞等待获取对方的锁。10.1 修复方案把 Promise 移出锁 原子状态更新修复提交给出的思路可以概括为Promise 提供的一些方法组合使用后具有类似compareAndSet的语义一种众所周知的、安全的模式。Permit.release()的新结构如下// old implementation // def release(): Unit self.synchronized { // val next waitq.pollFirst() // if (next ! null) next.setValue(this) // - nogo: still synchronized // else availablePermits 1 // } // new implementation // 这里定义一个专用锁对象而不是使用 this 引用本身。 // 它与内部 Queue 同义因为正是队列驱动了我们的同步需求 // 但使用这个专用引用给将来重构留下更易操作的空间。 private[this] final def lock: Object waitq tailrec def release(): Unit { // 把 Promise 传出锁外 val waiter lock.synchronized { val next waitq.pollFirst() if (next null) { availablePermits 1 } next } if (waiter ! null) { // 由于不再与中断处理器同步 // 我们借助 Promise 的原子状态在竞态时做出正确行为 if (!waiter.updateIfEmpty(Return(this))) { release() } }新实现更复杂且不再与中断处理器同步这暴露给开发者一个新的考量点中断 Promise 与把 Permit 交给它之间的任何竞态。updateIfEmpty以原子方式确保要么设置成功、要么重试从而既避免了在锁内执行任意 continuation又保持了正确性。10.2 哪些同步操作是低风险的把同步与 Promise 一起使用需要非常小心才能确保程序没有死锁风险。低风险的同步用法包括变更或读取一个字段在私有ArrayDeque上 push / pop 元素而同步块内的高风险动作包括调用由调用方注入的函数该函数可能阻塞、获取自己的锁、把 pi 计算到十亿位等等调用来源未知 trait 的方法本质上与调用用户注入的函数相同完成一个 Promise上面就是带危险 continuation 的示例十一、从失败中恢复rescue复合 Future 会在任一组成 Future 失败时失败。但常常需要从这类失败中恢复。Future上的rescue组合子是flatMap的对偶flatMap作用于值rescue作用于异常其余行为完全相同。通常我们希望只处理一部分可能的异常为此rescue接受一个PartialFunction把Throwable映射到Futuretrait Future[A] { .. def rescueB : A: Future[B] .. }下面的代码在请求因TimeoutException失败时无限重试import com.twitter.util.Future import com.twitter.finagle.http import com.twitter.finagle.TimeoutException def fetchUrl(url: String): Future[http.Response] ??? def fetchUrlWithRetry(url: String): Future[http.Response] fetchUrl(url).rescue { case exc: TimeoutException fetchUrlWithRetry(url) }由于PartialFunction只匹配声明过的异常类型其他异常如连接拒绝不会被拦截会按原样向上传播。十二、竞速Future.select 与 selectIndex有时候我们并不关心哪个 Future 先完成。三种常见场景备份或对冲请求backup / hedged requests发出两个相同请求寄希望于其中一个慢时另一个不慢。这通常是使用备份请求的好时机——Finagle 的MethodBuilder通过idempotent方法的maxExtraLoad参数提供了开箱即用的备份请求支持MethodBuilder.rst阈值时间过后若未收到响应就发出第二个请求并通过基于maxExtraLoad的RetryBudget防止延迟突变时备份请求泛滥maxExtraLoad设为0.0即禁用。并发工作希望一项工作完成时立刻排队更多工作。并发工作每个 Future 满足后都需要处理且哪个先满足无所谓。在不便使用备份请求的场合Future.select往往是正确的工具。Future.select有三种模式同样是 Scala 原生、Java 可用Futures.selectFuture实例上的方法在 Java 中也能直接用。最简单的是Future实例上的方法返回最先完成的那个 Future——Future#select与Future#or行为完全相同import com.twitter.util.Future val original: Future[Tweet] ??? val hedged: Future[Tweet] ??? // Future#selectU : A: Future[U] val fasterTweet original.select(hedged)更强大的是Future.select与Future.selectIndex。Future.select接受一组 Future 集合返回一个 Future其中包含第一个被满足的 Future 的内容与其余 Future 的集合。拿到其余 Future 非常有用可以中断它们若不再需要、检查它们是否也已满足并立即处理无需让出、或者再次对剩余 Future 执行 select 并让出直到其中之一满足。处理Future.select的结果通常利用 Future 递归。下面是三种典型用法用法一中断其余 Futureimport com.twitter.util.Future import com.twitter.util.Try import java.util.concurrent.CancellationException val doWork: () Future[Tweet] ??? val tweets: Seq[Future[Tweet]] Seq.fill(10)(doWork) // Future.selectA: Future[(Try[A], Seq[Future[A]])] val first: Future[Tweet] Future.select(tweets).flatMap { case (first, rest) val cancelEx new CancellationException(lost the race) rest.foreach { f f.raise(cancelEx) } Future.const(first) }用法二尽可能急切地处理已完成的import com.twitter.util.Future import com.twitter.util.Try val doWork: () Future[Tweet] val tweets: Seq[Future[Tweet]] Seq.fill(10)(doWork) def tweetSentiment(tweet: Tweet): Int ??? def aggregateTweetSentiment(f: Future[(Try[Tweet], Seq[Future[Tweet]])]): Future[Seq[Int]] f match { case (first, rest) val (finished, unfinished) (Future.const(first) : rest).foldLeft((Seq[Tweet](), Seq[Future[Tweet]]())) { case ((complete, incomplete), f) f.poll match { case Some(Return(tweet)) (complete : tweet, incomplete) case None (complete, incomplete : f) case _ (complete, incomplete) // failed future } } val sentiments finished.map(tweetSentiment) if (unfinished.isEmpty) Future.value(sentiments) else Future.select(unfinished).flatMap(aggregateTweetSentiment).map(sentiments _) } // Future.selectA: Future[(Try[A], Seq[Future[A]])] val avgTweetSentiment: Future[Int] Future.select(tweets).flatMap(aggregateTweetSentiment).map { seq if (seq.isEmpty) 0 else (seq.sum / seq.length) }这个递归组合子在每次一轮 select 后把已经满足f.poll返回Some(Return(...))的结果就地处理掉对未完成的继续 select实现随完成随处理的流水线效果。用法三持续竞速直到第一个成功结果import com.twitter.util.Future import com.twitter.util.Try val doWork: () Future[Tweet] val tweets: Seq[Future[Tweet]] Seq.fill(10)(doWork) def raceTheTweets(f: Future[(Try[Tweet], Seq[Future[Tweet]])]): Future[Tweet] f match { case (Throw(_), rest) if rest.length 1 // 有失败但还有更多 Future 可等 Future.select(rest).flatMap(raceTheTweets _) case (Throw(_), Seq(last)) // 只剩一个无论成败都返回它 last case (result, _) // 要么成功要么最后一个 Future 也失败了 Future.const(result) } // Future.selectA: Future[(Try[A], Seq[Future[A]])] val first: Future[Tweet] Future.select(tweets).flatMap(raceTheTweets _)Future.selectIndex一个更强大但略微笨重的 API。它从IndexedSeq[Future]中只返回哪个最先被满足返回下标Int。好处有二其一如果你对传入集合的序列顺序有额外信息可以据此做决策——相比之下Future.select返回的是哪个 Future 你是不知道的其二它避免了返回复杂类型只返回数组下标import com.twitter.util.Future import com.twitter.util.Try import java.util.concurrent.CancellationException val doWork: () Future[Tweet] val tweets: IndexedSeq[Future[Tweet]] IndexedSeq.fill(10)(doWork) // Future.selectIndexA: Future[Int] val first: Future[Tweet] Future.selectIndex(tweets).flatMap { idx val cancelEx new CancellationException(lost the race) for (i 0 until 10 if i ! idx) tweets(i).raise(cancelEx) tweets(idx) }十三、深入理解Twitter Futures 的调度与实现为了用好上述组合子有必要理解 Twitter Futures 的调度模型详见 developers/Futures.rst三种调度哲学Scala Future 默认让别人跑投递到线程池Java CompletableFuture 默认我现在就跑调用线程内联执行栈不安全Twitter Future 默认我稍后跑仍利于缓存且栈安全代价是栈轨迹不那么有用。Future.value(doWorkA()).map {...}.map {...}在 Twitter 模型下退化为在同一调用栈内顺序执行三个闭包从而获得性能与栈安全。中断InterruptionTwitter Future 允许设置中断处理器并能沿flatMap链向上传播中断——当 Future 不再需要时可以取消产生该 Future 的底层工作且默认行为正确、无需为每个flatMap特判。中断是建议性的可以选择忽略。内部实现主要有ConstFuture已满足的 Future本质是Try的薄包装与Promise真实世界中的大多数 Future先未满足、后被满足。Promise使用状态机Waiting、Interruptible、Transforming、Interrupted、Linked、Done。Promise#updateIfEmpty及其派生方法如Promise#setValue会在未满足时将 Promise 从任意状态推向 Done。Linked 状态允许无限 Future 递归而不泄漏空间。可分离 PromiseDetachable当多个执行依赖共享同一个 Future 的结果时中断它可能误伤他人但永不中断又可能在 Future 永不满足时造成内存泄漏。Promise.attached(underlying)用于显式构造以后可能不再需要的 Promise。这些实现细节正是上文 AsyncSemaphore 修复中updateIfEmpty语义的底层来源也是rescue、select递归组合子能够安全工作的基础。十四、相关学习资源Effective Scala中有一个专门讨论 futures 的章节给出了 Twitter 内部使用 Future 的风格建议。自 Scala 2.10 起Scala 标准库有了自己的 Future 实现与 API与 Finagle 使用的com.twitter.util.Future大体相似但存在命名差异如flatMap一致而调度与取消语义不同见上文对比。Akka文档中也有专门介绍 futures 的章节。Finagle Block Party详细解释了为什么阻塞是有害的更重要的是如何发现并修复阻塞对应上文blocking_ms指标。小结Finagle 的并发模型建立在com.twitter.util.Future之上Future 是廉价的轻量级线程通过flatMap让出控制权阻塞工作一律交给FuturePool承载Finagle 内部 DNS 解析、OffloadFilter全局卸载都基于此顺序组合用flatMap、并发组合用collect/collectToTry、并行组合用join的四种模式、失败恢复用rescue、竞速用select/selectIndex而任何在同步块内完成 Promise 或调用注入函数的操作都可能引入死锁需要像 AsyncSemaphore 修复那样把 Promise 移出锁并借助updateIfEmpty的原子语义来规避。掌握这组工具就能在 Finagle 中写出高并发、可组合且无阻塞的 RPC 代码。赞分享后端RPC框架【免费下载链接】finagleA fault tolerant, protocol-agnostic RPC system项目地址https://gitcode.com/gh_mirrors/fi/finagle点击查看免费下载相关推荐Comprehensive Rust futures组合join与select操作Comprehensive Rust futures组合join与select操作 在Rust异步编程中Future未来代表一个异步操作的结果而asy文档教程RustTraining 异步编程实战从零手写 TimerFuture、Join 与 Select 组合器RustTraining 异步编程实战从零手写 TimerFuture、Join 与 Select 组合器 导读 本文聚焦 RustTraining 仓库文档教程Rust 异步控制流实战async 通道、Join 与 Select 组合并发逻辑Rust 异步控制流实战async 通道、Join 与 Select 组合并发逻辑 本篇技术指南以 Google Comprehensive Rust 课程文档教程上一篇5分钟掌握Unity游戏模组制作MelonLoader终极配置与使用指南下一篇MelonLoader深度解析Unity游戏模组加载的终极高效方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表