news 2026/10/12 2:16:57

concurrent-ruby Cancellation 协作用户指南:用 Promise 事件实现优雅的任务取消、超时与并行协同

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
concurrent-ruby Cancellation 协作用户指南:用 Promise 事件实现优雅的任务取消、超时与并行协同
  • 开发工具

【免费下载链接】concurrent-ruby

Modern concurrency tools including agents, futures, promises, thread pools, supervisors, and more. Inspired by Erlang, Clojure, Scala, Go, Java, JavaScript, and classic concurrency patterns.

项目地址:https://gitcode.com/gh_mirrors/co/concurrent-ruby
点击查看免费下载

Cancellation 是 concurrent-ruby 提供的协作式(cooperative)取消抽象:它不直接杀线程,而是通过一个可被解析的 Promise 事件对象作为"取消源头",由任务本身周期性检查取消状态并自行决定如何退出。本文基于 docs-source/cancellation.out.md 及其配套示例 docs-source/cancellation.init.rb,完整演示取消、超时、任务链回调与取消组合的全部写法,并结合源码与测试说明其底层原理,读完即可在真实异步任务中落地使用。

为什么需要协作式取消

Ruby 内置的Thread#raise、Thread#kill以及Timeout模块虽然能强制终止线程,但可能在任意指令处中断代码,导致锁未释放、资源未清理、状态被破坏,被社区视为危险 API(concurrent-ruby 源码注释中亦引用了相关讨论,见 取消实现)。Cancellation正是为此提供的替代方案:

  • 取消是协作式的:任务持有取消对象,周期性检查自己是否被取消,然后以受控方式退出;
  • 取消通过解析一个 Promise 事件触发:origin(源头)被解析的瞬间,所有持有该取消对象的任务都会观测到canceled?变为true;
  • 取消可被组合、可被派生:一个取消可以 join 多个取消,也可以作为新的取消源头继续传递。

快速上手:让后台任务"跑直到被取消"

创建取消对象有两种角色(源码):

cancellation, origin = Concurrent::Cancellation.new # => #<Concurrent::Cancellation:0x000002 pending>

这里用到了 Ruby 的多重赋值:Cancellation#to_ary返回[self, @Origin],于是cancellation是取消对象本身(传给任务做检查),origin是Concurrent::Promises.resolvable_event(由调用方调用resolve来触发取消)。在 初始化实现 中,initialize(origin = Promises.resolvable_event)说明不传参数时默认自动创建一个可解析事件。

把取消对象传入异步任务,让任务循环检查canceled?:

require 'concurrent-edge' def do_stuff(*args) sleep 0.01 :stuff end cancellation, origin = Concurrent::Cancellation.new async_task = Concurrent::Promises.future(cancellation) do |cancellation| # 反复干活,直到被取消 do_stuff until cancellation.canceled? :stopped_gracefully end # => #<Concurrent::Promises::Future:0x000003 pending> sleep 0.01 # 稍等片刻,再通过解析 origin 停掉线程 origin.resolve # => #<Concurrent::Promises::ResolvableEvent:0x000004 resolved> async_task.value! # => :stopped_gracefully

任务通过until cancellation.canceled?观察取消状态,循环自然结束后返回:stopped_gracefully,Future#value!拿到正常结果——这就是"优雅停止":不做任何强制中断,任务自行走完退出路径。

注意:Concurrent::Cancellation属于 edge(实验性)命名空间,需要require 'concurrent-edge'引入,对应文件 concurrent-edge.rb 中的加载项。

取消时抛错:check! 与 CancelledOperationError

如果任务无法自行判断退出时机、希望取消时直接中断当前流程并让 Future 以错误终结,可以改用check!:

cancellation, origin = Concurrent::Cancellation.new async_task = Concurrent::Promises.future(cancellation) do |cancellation| while true cancellation.check! do_stuff end end sleep 0.01 origin.resolve async_task.result # => [false, # nil, # #<Concurrent::CancelledOperationError: Concurrent::CancelledOperationError>]

check!在 源码 中实现为:

def check!(error = CancelledOperationError) raise error if canceled? self end

即"已取消则抛错,否则返回自身",且允许自定义异常类型。CancelledOperationError定义于 errors.rb,继承自Concurrent::Error。此时Future#result返回三元组[fulfilled?, value, reason],即[false, nil, CancelledOperationError],任务以拒绝(rejected)状态结束。

取消后继续执行:在 origin 上挂链

取消不仅可以停止任务,还可以作为"事后回调"的触发器——例如记录取消日志、规划重新执行:

cancellation, origin = Concurrent::Cancellation.new cancellation.origin.chain do # 该块会在 Cancellation 被取消后执行 # 可用于记录取消原因,或安排新的重试任务 end # => #<Concurrent::Promises::Future:0x000008 pending>

origin本身是ResolvableEvent,chain返回一个在其解析后执行的 Future(实现见 chain_on)。这样取消流程就被纳入 Promise 链,可以继续组合then、rescue、zip等操作,实现"取消即调度"的重试策略。

限时执行:Cancellation 作为 Timeout 的替代

协作式取消还能实现安全的超时:不再用Timeout强行中断,而是预先安排一个在指定时间后被解析的 Promise 作为 origin。

timeout = Concurrent::Cancellation.new Concurrent::Promises.schedule(0.02) # => #<Concurrent::Cancellation:0x000009 pending> # 或使用快捷方法 timeout = Concurrent::Cancellation.timeout 0.02 # => #<Concurrent::Cancellation:0x00000a pending> count = Concurrent::AtomicFixnum.new Concurrent.global_io_executor.post(timeout) do |timeout| # 干活直到被取消 count.increment until timeout.canceled? end timeout.origin.wait # => #<Concurrent::Promises::Event:0x00000c resolved> count.value # => 177576

两种写法等价:Cancellation.timeout(intended_time)在 源码 中就是new Concurrent::Promises.schedule(intended_time)。Promises.schedule创建的是 ScheduledPromise,到点后自动解析其 event,进而触发取消。

超时语义由任务自己掌控:这里count.increment until timeout.canceled?,origin 解析后循环立即停止;由于没有强制中断,任务一定停在安全边界。示例运行中count.value约为 177576,说明在这 0.02 秒内完成了约 17.7 万次自增——该数值与运行环境相关,仅作参考。Concurrent.global_io_executor是全局 IO 执行器,定义于 configuration.rb。

并行任务共享一个取消:一个失败、全体取消

多个后台任务可以共享同一个取消对象:当其中一个失败时,由它解析 origin,其余任务随后观测到取消并优雅退出。示例让 4 个任务各自尝试计数到 100,但存在约 5% 概率的随机失败:

cancellation, origin = Concurrent::Cancellation.new tasks = 4.times.map do |i| Concurrent::Promises.future(cancellation, origin, i) do |cancellation, origin, i| count = 0 100.times do break count = :cancelled if cancellation.canceled? count += 1 sleep 0.001 if rand > 0.95 origin.resolve # 取消 raise 'random error' end count end end end Concurrent::Promises.zip(*tasks).result # => [false, # [:cancelled, nil, :cancelled, :cancelled], # [nil, #<RuntimeError: random error>, nil, nil]]

解读zip(...).result返回的三元组:

  • 第一个元素false:聚合 Future 整体未 fulfilled(因为某个任务抛错);
  • 第二个元素是各任务返回值数组:失败任务返回nil,其余返回:cancelled——它们都在循环中检测到取消后通过break count = :cancelled优雅退出;
  • 第三个元素是各任务 reason 数组:只有失败任务携带RuntimeError: random error。

注释掉随机失败部分后,所有任务都完成 100 次计数:

cancellation, origin = Concurrent::Cancellation.new tasks = 4.times.map do |i| Concurrent::Promises.future(cancellation, origin, i) do |cancellation, origin, i| count = 0 100.times do break count = :cancelled if cancellation.canceled? count += 1 sleep 0.001 count end end end Concurrent::Promises.zip(*tasks).result # => [true, [100, 100, 100, 100], nil]

即[true, [100, 100, 100, 100], nil]:全部 fulfilled,无异常。这个模式正是"故障快速扩散"的典型场景:任一任务出错,立即取消全体,避免其他任务继续空转。zip 的聚合语义见 zip_futures_on。

组合取消:join 与手动 AND 语义

join:任一取消即取消

join产生的新取消对象,在自身或任一被合并取消对象被取消时触发:

cancellation_a, origin_a = Concurrent::Cancellation.new cancellation_b, origin_b = Concurrent::Cancellation.new combined_cancellation = cancellation_a.join(cancellation_b) # => #<Concurrent::Cancellation:0x000019 pending> origin_a.resolve # => #<Concurrent::Promises::ResolvableEvent:0x00001a resolved> cancellation_a.canceled? # => true cancellation_b.canceled? # => false combined_cancellation.canceled? # => true

其实现(源码)正是"任意一个解析即解析"的any_event语义:

def join(*cancellations) Cancellation.new Promises.any_event(*[@Origin, *cancellations.map(&:origin)]) end

Promises.any_event在 promises.rb 中描述为"第一个 futures/events 解析后即解析"。any_event内部使用 AnyResolvedEventPromise,只要任一阻塞源解析就立即解析。

手动组合:AND 语义(全部取消才取消)

如果业务要求"两个取消条件同时满足才取消",可绕过join直接构造:把两个 origin 用&(即zip)组合后作为新 Cancellation 的 origin。Event 的&别名指向zip(见 promises.rb),语义为"两者都解析后才解析":

cancellation_a, origin_a = Concurrent::Cancellation.new cancellation_b, origin_b = Concurrent::Cancellation.new # 仅当 a 与 b 都被取消时才取消 combined_cancellation = Concurrent::Cancellation.new origin_a & origin_b # => #<Concurrent::Cancellation:0x00001d pending> origin_a.resolve # => #<Concurrent::Promises::ResolvableEvent:0x00001e resolved> cancellation_a.canceled? #=> true cancellation_b.canceled? #=> false combined_cancellation.canceled? #=> false origin_b.resolve # => #<Concurrent::Promises::ResolvableEvent:0x00001f resolved> combined_cancellation.canceled? #=> true

可以看到,只有origin_a和origin_b都被解析后,combined_cancellation才进入 canceled 状态。这也体现了设计上的灵活性:Cancellation#initialize接受任意Promises::Future/Promises::Event作为 origin(源码),因此可以把任意复杂的 Promise 组合表达式当作取消源头。

底层原理:canceled? 只是 origin 的 resolved?

取消状态本身没有任何独立状态机,canceled?的实现极其简洁(源码):

def canceled? @Origin.resolved? end

也就是说:"是否被取消" 完全等价于 "origin 是否已被解析"。取消机制的可靠性因此完全由 Promises 框架的状态管理保证:Event/Future 一旦解析即不可逆,所有依赖它的链、回调与检查点会同时收到通知(内部状态Pending/Resolved的判定见 promises.rb,resolved?在 promises.rb)。to_s/inspect则按此状态输出pending或canceled(源码)。

这种设计带来的好处:

  • 取消状态天然线程安全,无需额外加锁;
  • origin 可以是任何 Future/Event,取消触发条件可以极其复杂(超时、AND、OR、链式依赖等);
  • 取消与 Promise 体系无缝衔接,可直接在 origin 上chain、then、rescue。

进阶:用 ResolvableFuture 承载取消原因

origin 不限于 Event,也可以是ResolvableFuture——此时取消可以携带值或错误原因。参考 cancellation_spec.rb 中的用法:

cancellation, origin = Concurrent::Cancellation.new(Concurrent::Promises.resolvable_future) origin.resolve false, nil, err = StandardError.new('Cancelled') expect(cancellation.canceled?).to be_truthy cancellable_branch = Concurrent::Promises.delay { 1 } expect((cancellable_branch | origin.to_future).reason).to eq err

Promises.resolvable_future创建可解析的 Future(源码),resolve(false, nil, err)以拒绝方式解析它,canceled?立即为真,同时错误err可沿 Promise 链传播,供下游任务分析取消原因。|是any的别名(promises.rb),delay制造惰性分支(promises.rb)。这为"取消带原因、带数据"的高级场景提供了扩展点。

测试验证

仓库在 spec/concurrent/cancellation_spec.rb 中对本文所有核心行为给出了断言:

  • 基础行为(spec 第 5-33 行):两个任务分别用canceled?循环和check!循环等待取消,origin.resolve后前者正常返回:done,后者以CancelledOperationError结束,同时cancellation.to_s从pending变为canceled;
  • 取消即触发链(spec 第 35-43 行):origin 解析后canceled?为真,且可参与|(any)组合;
  • 带原因的取消(spec 第 53-61 行):以resolvable_future为 origin,resolve false, nil, err后取消成立且 reason 可传播;
  • join 语义(spec 第 81-91 行):cancellation_a.join(cancellation_b)在origin_a.resolve后立即取消,与文档示例完全一致。

这些测试同时印证了文档示例的可运行性:文档配套的 cancellation.init.rb 中do_stuff以sleep 0.01模拟耗时工作,测试则用Thread.pass让出时间片,两者验证的是同一套 API 语义。

实践建议

  • 优先协作式取消:凡是长循环、轮询、批处理任务,都应传递cancellation并周期性检查canceled?/check!,避免使用Thread#kill等强杀手段;
  • 阻塞操作配超时:源码注释建议(cancellation.rb)——所有阻塞动作以带超时的方式循环执行,超时后先检查取消状态,未被取消再继续阻塞,这样协作式取消不会因长期阻塞而失效;
  • 取消后必做清理:在循环退出路径中统一处理资源释放、状态回写,再返回结果或抛出CancelledOperationError;
  • 组合取消按需选择:join(OR 语义)适合"任一条件满足即取消"(如超时或用户中断);手动&(AND 语义)适合"多条件齐备才取消"(如等待多个前置任务全部结束);
  • 注意 edge 标记:Concurrent::Cancellation位于 edge 命名空间(edge.rb 说明其为实验性、API 可能变动的功能),生产使用前应确认当前版本 API 并锁定依赖版本。

总结

concurrent-ruby 的Cancellation用最轻量的方式解决了并发编程中最棘手的"安全停止"问题:取消只是一个可解析的 Promise 事件,任务端只需周期性检查,触发端只需origin.resolve。配合schedule可实现安全超时,配合join/&可实现 OR/AND 组合取消,配合chain/rescue可编排取消后的恢复流程。文档配套示例 docs-source/cancellation.out.md、初始化代码 docs-source/cancellation.init.rb 与测试 spec/concurrent/cancellation_spec.rb 三者相互印证,可直接作为实战蓝本。

  • 开发工具

【免费下载链接】concurrent-ruby

Modern concurrency tools including agents, futures, promises, thread pools, supervisors, and more. Inspired by Erlang, Clojure, Scala, Go, Java, JavaScript, and classic concurrency patterns.

项目地址:https://gitcode.com/gh_mirrors/co/concurrent-ruby
点击查看免费下载
上一篇:Bazel 远程缓存完全指南:协议、后端搭建与实战配置
下一篇:LeetCode-Go 题解精讲:1143. Longest Common Subsequence(最长公共子序列,二维动态规划)

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/12 2:16:27

数据库课程设计实战:教务管理系统从E-R图到SQL Server建库全流程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/12 2:15:22

CLion+CMake+OpenOCD:构建树莓派Pico高效调试环境

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华