concurrent-ruby Cancellation 协作用户指南:用 Promise 事件实现优雅的任务取消、超时与并行协同
开发工具【免费下载链接】concurrent-rubyModern 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模块虽然能强制终止线程但可能在任意指令处中断代码导致锁未释放、资源未清理、状态被破坏被社区视为危险 APIconcurrent-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_gracefullyFuture#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 pendingorigin本身是ResolvableEventchain返回一个在其解析后执行的 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.rbresolved?在 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.rbdelay制造惰性分支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为 originresolve 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组合取消按需选择joinOR 语义适合任一条件满足即取消如超时或用户中断手动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-rubyModern 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点击查看免费下载相关推荐concurrent-ruby 协作式取消Cancellation实战从基础轮询到超时与多任务组合concurrent ruby 协作式取消Cancellation实战从基础轮询到超时与多任务组合 Cancellation 是 concurrent r开发工具Concurrent Ruby Cancellation机制安全取消任务的完整解决方案Concurrent Ruby Cancellation机制安全取消任务的完整解决方案 在现代并发编程中安全地取消运行中的任务是一个关键挑战。Concurr开发工具深入理解 Rust 异步任务取消Cancellation超时、取消点与优雅清理实战解析深入理解 Rust 异步任务取消Cancellation超时、取消点与优雅清理实战解析 导读 本篇文章围绕 100 exercises to learn示例工程教程上一篇Bazel 远程缓存完全指南协议、后端搭建与实战配置下一篇LeetCode-Go 题解精讲1143. Longest Common Subsequence最长公共子序列二维动态规划创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考