- 开发工具
【免费下载链接】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.
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)]) endPromises.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 errPromises.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.
相关推荐
concurrent-ruby 协作式取消(Cancellation)实战:从基础轮询到超时与多任务组合
concurrent ruby 协作式取消(Cancellation)实战:从基础轮询到超时与多任务组合 Cancellation 是 concurrent r
开发工具Concurrent Ruby Cancellation机制:安全取消任务的完整解决方案
Concurrent Ruby Cancellation机制:安全取消任务的完整解决方案 在现代并发编程中,安全地取消运行中的任务是一个关键挑战。Concurr
开发工具深入理解 Rust 异步任务取消(Cancellation):超时、取消点与优雅清理实战解析
深入理解 Rust 异步任务取消(Cancellation):超时、取消点与优雅清理实战解析 导读 本篇文章围绕 100 exercises to learn
示例工程教程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考