Java常用类库-16-ListenableFuture的执行线程与取消边界
结果完成以后,线程还在做什么
商品查询完成后需要转换价格、记录结果,再通知调用方。把这些操作注册成 Future 回调,能够避免在查询入口直接阻塞等待,但并不能据此断定回调运行在另一个线程。使用 directExecutor 时,完成查询的线程可能继续执行全部回调;若回调等待慢速存储,查询线程就仍然被占用。
ListenableFuture 在普通 Future 上增加完成监听器,每个监听器关联一个 Executor。这个关联才决定执行策略。“异步结果”描述何时取得值,Executor 描述由谁执行动作,两者应分别建模。本章固定 Guava 33.5.0-jre,以 Java 8 代码在 Zulu 8u472 与 Corretto 21.0.11 运行。ListenableFuture 文档
三个时间点必须分开
查询计算完成、Future 进入终态、监听器执行完成,是三个不同时间点。Future 终态包括正常值、失败和取消;isDone 不能证明成功,也不能证明所有监听器已完成。若后续业务依赖回调产生的数据,应返回包含该阶段的派生 Future,或提供专门的完成信号。
本章先创建未完成的 SettableFuture,再注册一个直接执行的监听器。名为 catalog-setter 的线程调用 set("SKU-1");监听器记录线程名后进入 latch 等待。此时测试线程观察到 source.isDone() 已为 true,但 set 所在线程尚未返回。释放 latch 后,监听器结束,set 才返回。
1 | |
这个实验证明了直接监听器能阻塞完成线程,也说明“Future 已完成”和“完成操作的调用已经返回”并非同一个条件。调用方若只等待 get(),仍不应读取由任意附加监听器写入、却没有独立同步关系的共享状态。
再对已经完成的 Future 注册 direct listener,执行线程变成注册监听器的测试线程。其原因是无需等待未来的完成事件,注册路径可以直接安排执行。directExecutor 不创建线程,也不会承诺固定线程名;它在调用 execute 的线程内运行 Runnable。MoreExecutors.directExecutor
成功、失败与转换使用同一执行策略
需要获得结果值时,Futures.addCallback 比原始 addListener 更直接:正常完成进入 onSuccess,失败进入 onFailure。转换结果则可使用 Futures.transform,它返回一个新的 ListenableFuture。把转换注册到命名线程池后,测试得到商品编号长度 5,并记录实际执行线程为 catalog-callback。
失败实验使用已经完成的失败 Future,原因是 IllegalArgumentException("bad SKU")。回调仍通过指定的线程池执行,onFailure 得到失败原因。测试用 latch 等待回调,再断言消息,避免把“回调已注册”误当成“回调已经调用”。普通失败和执行线程选择是两项互不替代的验收。
以下完整 Java 8 示例把转换安排到独立执行器,等待的是转换完成后的 Future。它在结束时关闭线程池,并明确等待终止;若线程池无法结束,就抛出错误而非让进程悄悄留下线程。
1 | |
transform 适合函数立即得到结果;如果函数返回另一个异步结果,应选择表达异步展开的 API,例如 transformAsync,而不是留下 Future<Future<T>> 后让调用者猜测需要等哪一层。本章没有把异步展开纳入执行数据,因此结论聚焦于同步转换与明确的 Executor。
拒绝执行发生在哪一层
一个有界线程池可能拒绝回调;一个已经关闭的线程池也不能保证继续接受任务。因此,成功产生源值并不保证转换还能执行。实验为 Futures.transform 提供必定抛出 RejectedExecutionException("full") 的 Executor,最终派生 Future 的 get 抛 ExecutionException,原因是拒绝异常。
这条行为来自固定版本 AbstractTransformFuture.create:注册输出 Future 时使用 rejectionPropagatingExecutor,把提交阶段的拒绝反映为输出失败。不能据此宣称所有 addListener 或 addCallback 都会把拒绝写回原始 Future;普通监听器是附加动作,原始计算可能已经成功完成。AbstractTransformFuture 源码
AbstractFuture.executeListener 对 Executor 执行中的 RuntimeException 采用记录日志的处理方式。选择派生 Future 或独立监听器,决定了调用者有没有一个能等待并观察失败的结果对象。导出文件、持久化确认等必须成功的阶段,不能只挂一个无人观察的回调。AbstractFuture 源码
无界线程池虽然减少了队列满时的拒绝,却可能把压力转移为线程或内存增长。这里没有用“不会拒绝”作为执行器选型目标。真正的约束应包括队列容量、并行数、超时、调用方如何处理拒绝,以及是否允许部分结果丢失。测试只注入确定拒绝,不模拟生产容量或压测吞吐。
取消包含状态变化与执行控制
对 transform 返回的 Future 调用 cancel,固定版本会尝试向输入 Future 传播取消。实验取消派生结果后确认源 SettableFuture 也处于 cancelled 状态。SettableFuture 本身没有正在运行的业务线程,因此这项断言只证明状态传播,不能证明线程已中断。
另一个实验通过 ListeningExecutorService 提交任务,让任务在线程中等待 latch。任务确实开始后调用 cancel(true),阻塞点收到 InterruptedException,测试才判定中断已送达。任务采用可中断等待,是实验成立的必要条件;如果业务吞掉中断或阻塞在不响应中断的操作中,Future 取消不等于任务必定及时停止。
CompletableFuture 的 cancel 参数有不同契约:它明确说明 mayInterruptIfRunning 不影响该实现的处理。对已经开始的 supplyAsync supplier 取消后,测试主动释放 latch,让 supplier 记录中断标记,观察到 false。取消的是结果状态,supplier 仍执行到自己的结束位置。CompletableFuture.cancel
取消效果由Future类型、任务提交方式和中断响应共同决定。JDK ExecutorService返回的普通Future也提供取消机制。实验只比较ListeningExecutorService任务与CompletableFuture supplier这两条具体路径,不能扩展成整个Guava与JDK的线程停止能力排名。
一个线程池同时承载查询与阻塞回调,还可能产生饥饿:全部工作线程等待仍排在同一队列里的任务,就没有线程推进结果。简单地把 directExecutor 换成任意线程池,并不能证明系统不会阻塞。应沿依赖图检查谁等待谁,尤其避免在同一个受限执行器内同步 get 后续阶段。
共享状态和资源生命周期
监听器的注册与开始执行有文档规定的内存一致性关系,但这不能扩展成监听器结束与任意第三方读取之间自动有序。实验使用 AtomicReference 保存线程名和失败原因,再通过 latch 或 Future.get 完成观察;这些同步手段都属于断言能够可靠读取结果的条件。
本文没有假定监听器按注册顺序执行。ListenableFuture 不保证监听器执行顺序,使用多线程执行器时完成顺序还取决于调度。需要严格的业务先后关系时,应把后续阶段建立在前一阶段的派生结果上,而不能依赖两次 addListener 的书写次序。
资源所有权也需要明确。示例创建执行器,因此在 finally 中关闭;如果执行器由应用容器共享,单个查询回调就不应随意 shutdown。测试中的 latch 在 finally 释放,线程 join 与 awaitTermination 都有超时,避免断言失败后留下阻塞任务。对真实服务而言,关闭过程还需决定未完成结果如何交付调用者。
Future 的超时等待也不会自动等同于取消任务。调用者 get 超时后可以选择取消、放弃等待或转为后台跟踪;任何一种选择都应有明确资源策略。本文用五秒超时保护测试不永久挂起,而没有把这个测试超时解释为业务查询 SLA。
把执行策略放到业务接口上
价格转换若只做少量纯计算,直接执行可以减少任务排队;如果转换还访问库存、数据库或远程服务,就应把这些阶段单独建模。否则调用者看到的方法名仍是“转换价格”,实际却在查询线程上承担了另一段不可控等待。选择执行器时,应同时确定最大排队数量、拒绝后的结果状态,以及负责关闭线程池的组件。
线程池隔离也不能消除所有依赖环。假设一个单线程执行器正在执行回调,回调提交第二个任务到同一个执行器,然后调用 get 等待第二个任务。第二个任务排在当前回调之后,当前回调又等它完成,就无法前进。把 directExecutor 改成线程池只是改变了等待位置,仍需检查任务之间的依赖。此类转换应继续返回组合后的 Future,让执行线程退出当前阶段。
结果链和旁路监控也应区分。决定是否接受订单的库存结果必须成为调用者等待的结果链;记录一个非关键计数可以是独立监听器,但必须有独立的失败观察方式。若旁路动作实际上要求可靠送达,仅注册进程内回调并不提供持久化保证,进程退出前未完成的动作仍可能丢失。
迁移现有 Future 接口时,可先列出调用方当前依赖的四项行为:失败包装类型、回调线程、取消向上游传播的范围,以及关闭时仍在运行的任务如何收尾。再把这些行为写成有确定同步点的测试。只替换返回类型而不检查这些行为,编译通过也可能改变请求线程的占用时间和资源回收方式。本章的取消测试刻意等待任务已启动,避免把“尚未开始所以没有中断可观察”误写成实现不支持中断。
实测与练习
完整测试及运行说明覆盖四组场景。两个 JDK 均为 4 项测试、0 失败、0 错误、0 跳过;stdout 记录实际线程与时间线。这里的 latch 用于确定因果关系,没有以运行用时排名执行器。| 场景 | 观察结果 |
|---|---|
| 指定转换执行器 | 结果5,线程catalog-callback |
| 未完成源上的direct listener | catalog-setter被监听器阻塞,源已done |
| 已完成源上的direct listener | 注册线程执行 |
| transform执行器拒绝 | 派生结果失败,原因为RejectedExecutionException |
| 两种具体取消路径 | Guava阻塞任务收到中断;CompletableFuture supplier未收到 |
手算题:在一个已经完成的 Future 上注册 direct listener,listener 内等待另一个线程释放锁。哪个调用可能被阻塞?注册监听器的调用可能无法返回,因为执行器直接在该线程运行回调。Future 完成时间再早也不会改变这个执行策略。
改动练习:把直接监听器改成独立单线程执行器,保留 latch,分别断言源 set 返回与回调完成的顺序。再关闭执行器注入拒绝,观察普通监听器与 transform 输出 Future 的可观察差别,明确哪条结果链必须由调用者等待。
| 可迁移做法 | 适用边界 |
|---|---|
| 分开记录结果终态与回调完成 | 回调还有持久化、通知或其他副作用 |
| 按具体Future实现验收取消传播和中断 | 从一种异步API迁移到另一种API |
