☰
Finagle Futures 并发编程完全指南:从 FuturePool 到 flatMap / collect / join / select 的组合实战
2026/9/25 1:24:24 网站建设 项目流程
  • 后端
  • RPC框架

【免费下载链接】finagle

A 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/O(JDBC 驱动、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 组合。

Scala:

import com.twitter.util.{Future, FuturePool} def someIO(): String = // does some blocking I/O and returns a string val futureResult: Future[String] = FuturePool.unboundedPool { someIO() }

Java:

import com.twitter.util.Future; import com.twitter.util.FuturePools; import static com.twitter.util.Function.func0; Future<String> futureResult = FuturePools.unboundedPool().apply( func0(() -> someIO()); );

3.1 全局卸载(Offloading)

FuturePool 也是 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.auto=true

或者手动调优线程配置:

-com.twitter.finagle.offload.numWorkers=14 -com.twitter.finagle.netty4.numWorkers=10

对应的参数实现在 OffloadFuturePool.scala 与 numWorkers.scala 中:numWorkers未显式指定且auto开启时,会按com.twitter.jvm.numProcs().ceil推导;池的默认兜底是FuturePool.unboundedPool。卸载类过滤器的行为在 OffloadFilterTest.scala 中有大量测试覆盖。

开启全局卸载后,还可选启用实验性的 Offload 准入控制(根据工作队列等待时间拒绝工作,默认阈值 20ms):

-com.twitter.finagle.offload.admissionControl=enabled -com.twitter.finagle.offload.admissionControl=50.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 B,B 等待 A "用完" Person A——死锁。一个详细的真实死锁示例见下文"组合中的同步与死锁"。

六、Future 作为容器:三种状态与回调

被 Future 表示的常见操作包括:

  • 对远端主机的 RPC
  • 另一个线程中的长计算
  • 从磁盘读取

注意这些操作都可能失败:远端主机可能崩溃、计算可能抛异常、磁盘可能损坏。因此一个Future[T]恰好占据三种状态之一:

  • Empty(pending):尚未完成
  • 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 的流程),它通常包含三步:

  1. 抓取主页
  2. 解析页面找到第一个图片链接
  3. 抓取该图片链接

这是顺序组合的典型场景:要做下一步,必须先成功完成上一步。在 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.collect

Future 也可以并发地组合。扩展上面的例子:一次性抓取所有图片。并发组合由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.join);Future实例上的方法在 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 多个不同结果,有大量重载以支持不同数量的 Future:

import 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映射到Future:

trait 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 先完成。三种常见场景:

  1. 备份或对冲请求(backup / hedged requests):发出两个相同请求,寄希望于其中一个慢时另一个不慢。这通常是使用备份请求的好时机——Finagle 的MethodBuilder通过idempotent方法的maxExtraLoad参数提供了开箱即用的备份请求支持(MethodBuilder.rst),阈值时间过后若未收到响应就发出第二个请求,并通过基于maxExtraLoad的RetryBudget防止延迟突变时备份请求泛滥;maxExtraLoad设为0.0即禁用。
  2. 并发工作:希望一项工作完成时立刻排队更多工作。
  3. 并发工作:每个 Future 满足后都需要处理,且哪个先满足无所谓。

在不便使用备份请求的场合,Future.select往往是正确的工具。Future.select有三种模式,同样是 Scala 原生、Java 可用(Futures.select;Future实例上的方法在 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 递归。下面是三种典型用法:

用法一:中断其余 Future:

import 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 模型下退化为在同一调用栈内顺序执行三个闭包,从而获得性能与栈安全。
  • 中断(Interruption):Twitter 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 递归而不泄漏空间。
  • 可分离 Promise(Detachable):当多个执行依赖共享同一个 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框架

【免费下载链接】finagle

A fault tolerant, protocol-agnostic RPC system

项目地址:https://gitcode.com/gh_mirrors/fi/finagle
点击查看免费下载
上一篇:5分钟掌握Unity游戏模组制作:MelonLoader终极配置与使用指南
下一篇:MelonLoader深度解析:Unity游戏模组加载的终极高效方案

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询