
Ruby Fiber 与 Fiber::Scheduler 完全指南协作式并发、非阻塞 I/O 与自定义调度器实现【免费下载链接】rubyThe Ruby Programming Language项目地址: https://gitcode.com/GitHub_Trending/ru/ruby本文以 Ruby 官方文档 doc/language/fiber.md 为核心骨架结合 cont.c 源码与 test/fiber 下的真实测试用例系统讲解 Ruby Fiber 的协作式并发模型、Fiber::Scheduler调度器接口的完整设计以及如何在非阻塞执行上下文中使用Fiber.schedule、Fiber.set_scheduler等 API。读完本文你将掌握 Fiber 的上下文切换机制、调度器 14 个 hook 的语义与实现要点、IO#close中断阻塞 fiber 的底层时序并能参照仓库内的最小调度器实现写出自己的事件循环调度器。一、Fiber协作式并发的基石Fiber纤程为 Ruby 提供了一种协作式并发cooperative concurrency机制。与抢占式调度的线程不同Fiber 由程序自身显式让出控制权因此切换开销更小、行为更可预测。1.1 上下文切换yield、resume 与 transferFiber 执行用户提供的代码块。在块执行期间可以调用Fiber.yield或Fiber.transfer切换到其他 FiberFiber#resume则用于从上次Fiber.yield让出的位置继续执行。官方文档给出了最经典的流程控制示例#!/usr/bin/env ruby puts 1: Start program. f Fiber.new do puts 3: Entered fiber. Fiber.yield puts 5: Resumed fiber. end puts 2: Resume fiber first time. f.resume puts 4: Resume fiber second time. f.resume puts 6: Finished.运行这段程序输出顺序严格为1 → 2 → 3 → 4 → 5 → 6第一次resume进入 Fiber 执行到Fiber.yield控制权交还主执行流第二次resume从让出点继续打印 Resumed fiber. 后 Fiber 块自然结束控制权再次交还。这个简单的乒乓过程演示了 Fiber 的全部本质——显式的、可暂停与恢复的执行上下文。在 cont.c 中Fiber#resume的实现文档cont.c 中rb_fiber_resume相关注释进一步说明了语义resume从最后一次Fiber.yield的位置继续执行传给resume的参数会成为Fiber.yield表达式的返回值或块参数Fiber.yield传入的参数则会作为resume的返回值。借助这一双向传值机制Fiber 天然适合实现生成器generator、状态机与轻量协作任务。二、Fiber::Scheduler拦截阻塞操作的统一接口Fiber 本身只解决如何切换不解决何时切换。为此 Ruby 定义了Fiber::Scheduler接口它用于拦截阻塞操作如sleep、IO 读写、进程等待把让出控制权的决策权交给调度器。一个典型的实现是对EventMachine、Async这类事件循环库的封装。这种设计带来了清晰的关注点分离事件循环的实现细节与应用程序代码解耦同时支持分层调度器layered schedulers——多个调度器可以叠加用于埋点、监控等插桩instrumentation用途。2.1 设置与移除调度器为当前线程设置调度器只需一行Fiber.set_scheduler(MyScheduler.new)当线程退出时Ruby 会隐式调用Fiber.set_scheduler(nil)从源码看Fiber.set_scheduler的 C 实现cont.c#L2605-L2624文档明确说明设置调度器后非阻塞 Fiber通过Fiber.new(blocking: false)或Fiber.schedule创建在遇到可能阻塞的操作时会调用该调度器的 hook 方法线程终结时会调用调度器的close方法让调度器有机会妥善管理所有未完成的 Fiber。与调度器相关的三个查询方法语义各不相同cont.c#L2575-L2603Fiber.scheduler返回当前线程最后设置的调度器未设置时为nilFiber.current_scheduler仅当当前 Fiber 是非阻塞的时才返回调度器否则返回nil。test/fiber/test_scheduler.rb 中的test_current_scheduler验证了这一区别在主执行流中Fiber.scheduler有值而Fiber.current_scheduler为nil进入Fiber.schedule创建的 fiber 后Fiber.current_scheduler才返回调度器实例。2.2 设计理念无观点的轻量薄层调度器接口被刻意设计为无观点un-opinionated的轻量层介于用户代码与阻塞操作之间。核心约束是hook 不应翻译或转换参数与返回值——理想情况下用户代码传入的参数原封不动地交给调度器 hook返回值也原样返回。只有保持这种薄与直通调度器才能与各类上层框架无缝协作也才能被安全地叠加使用。2.3 需要实现的完整接口清单调度器是一个普通 Ruby 对象你可以自由实现其方法可选的 hook 会通过 Ruby 的响应性检测决定是否调用。以下是官方文档给出的完整接口骨架必须逐项理解其职责class Scheduler # 等待指定的进程 ID 退出。 # 此 hook 是可选的。 # parameter pid [Integer] 要等待的进程 ID。 # parameter flags [Integer] 适用于 Process::Status.wait 的标志位掩码。 # returns [Process::Status] 进程状态实例。 def process_wait(pid, flags) Thread.new do Process::Status.wait(pid, flags) end.value end # 在指定超时时间内等待给定 io 的可读性匹配指定事件。 # parameter event [Integer] IO::READABLE、IO::WRITABLE、IO::PRIORITY 的位掩码。 # parameter timeout [Numeric] 等待事件的时间秒。 # returns [Integer] 已就绪的事件的子集。 def io_wait(io, events, timeout) end # 从给定 io 读取数据到指定缓冲区。 # parameter io [IO] 要读取的 io。 # parameter buffer [IO::Buffer] 要写入的缓冲区。 # parameter offset [Integer] 缓冲区中的写入偏移。 # parameter length [Integer] 单次操作的最大读取量。 def io_read(io, buffer, offset, length) end # 从给定 io 的指定位置读取数据到指定缓冲区。 # parameter io [IO] 要读取的 io。 # parameter buffer [IO::Buffer] 要写入的缓冲区。 # parameter from [Integer] io 中的读取位置。 # parameter offset [Integer] 缓冲区中的写入偏移。 # parameter length [Integer] 单次操作的最大读取量。 def io_pread(io, buffer, from, offset, length) end # 从指定缓冲区写入数据到指定 IO。 # parameter io [IO] 要写入的 io。 # parameter buffer [IO::Buffer] 要读取的缓冲区。 # parameter offset [Integer] 缓冲区中的读取偏移。 # parameter length [Integer] 单次操作的最大写入量。 def io_write(io, buffer, offset, length) end # 从指定缓冲区写入数据到给定 io 的指定位置。 # parameter io [IO] 要写入的 io。 # parameter buffer [IO::Buffer] 要读取的缓冲区。 # parameter from [Integer] io 中的写入位置。 # parameter offset [Integer] 缓冲区中的读取偏移。 # parameter length [Integer] 单次操作的最大写入量。 def io_pwrite(io, buffer, from, offset, length) end # 让当前任务休眠指定时长未指定时长则永久休眠。 # parameter duration [Numeric] 休眠时间秒。 def kernel_sleep(duration nil) end # 执行给定块。若块执行超过指定超时时间则抛出指定的异常 klass。 # 通常只有进入调度器的非阻塞方法才会抛出此类异常。 # parameter duration [Integer] 等待时长超过后抛出异常。 # parameter klass [Class] 要抛出的异常类。 # parameter *arguments [Array] 传给异常构造函数的参数。 # yields {...} 要执行的用户代码。 def timeout_after(duration, klass, *arguments, block) end # 将主机名解析为 IP 地址数组。 # 此 hook 是可选的。 # parameter hostname [String] 示例www.ruby-lang.org。 # returns [Array] 主机名解析出的 IPv4 和/或 IPv6 地址字符串数组。 def address_resolve(hostname) end # 阻塞调用方 fiber。 # parameter blocker [Object] 等待的对象仅作信息用途。 # parameter timeout [Numeric | Nil] 等待时间秒。 # returns [Boolean] 阻塞操作是否成功。 def block(blocker, timeout nil) end # 解除对指定 fiber 的阻塞。 # parameter blocker [Object] 等待的对象仅作信息用途。 # parameter fiber [Fiber] 要解除阻塞的 fiber。 # reentrant 线程安全。 def unblock(blocker, fiber) end # 拦截非阻塞 fiber 的创建。 # returns [Fiber] def fiber(block) Fiber.new(blocking: false, block) end # 线程退出时调用。 def close self.run end def run # 在这里实现事件循环。 end end注意未来可能引入更多 hookRuby 将采用**特性检测feature detection**的方式按需启用这些新 hook——即通过检查调度器对象是否响应某个方法respond_to?来决定是否调用因此你的调度器无需为尚未使用的 hook 预留空实现。2.4 最小合法接口与强制约束从测试用例可以反推出 Ruby 真正强制的接口底线。test/fiber/test_scheduler.rb#L84-L107 的test_minimal_interface显示一个调度器至少需要实现block、unblock、io_wait、kernel_sleep四个方法以及fiber_interrupt。而test_fiber_interrupt_is_requiredtest/fiber/test_scheduler.rb#L109-L121进一步验证即使实现了block/unblock/io_wait/kernel_sleep若缺少fiber_interruptFiber.set_scheduler会抛出ArgumentError错误信息为Scheduler must implement #fiber_interrupt。这说明fiber_interrupt是当前版本调度器的强制性 hook与 doc/language/fiber.md 中IO#close一节描述的对应Fiber::Scheduler#fiber_interrupthook 是必需的完全一致。三、非阻塞执行让调度器接管阻塞点调度器 hook只会在特殊的非阻塞执行上下文non-blocking execution context中生效。需要强调的是非阻塞执行上下文会引入非确定性调度器 hook 的执行可能在程序中插入额外的上下文切换点程序的运行时序因此不再与源码书写顺序一一对应。3.1 创建非阻塞 Fiber用Fiber.new创建非阻塞上下文Fiber.new do puts Fiber.current.blocking? # false # 可能调用 Fiber.scheduler.io_wait。 io.read(...) # 可能调用 Fiber.scheduler.io_wait。 io.write(...) # 一定会调用 Fiber.scheduler.kernel_sleep。 sleep(n) end.resumeRuby 3.0 起引入了非阻塞 fiber概念cont.c#L2085-L2105非阻塞 fiber 遇到本会阻塞的操作如sleep、等待进程或 I/O时会把控制权让给其他 fiber由调度器负责阻塞管理与唤醒。前提有两个一是 fiber 以Fiber.new(blocking: false)创建默认值即为 false二是当前线程已通过Fiber.set_scheduler设置调度器。若线程未设置调度器阻塞与非阻塞 fiber 的行为完全相同。Ruby 还提供了一个简化创建非阻塞 fiber 的方法Fiber.scheduleFiber.schedule do puts Fiber.current.blocking? # false endFiber.schedule的 C 实现文档cont.c#L2528-L2568给出了预期行为立即在一个独立的非阻塞 fiber 中运行给定块首次遇到阻塞操作时让出控制权给外部执行流事件循环结束时由调度器恢复所有被阻塞的 fiber。文档特别提醒具体行为完全取决于当前调度器对Fiber::Scheduler#fiber的实现Ruby 并不强制Fiber.schedule的特定行为若未设置调度器调用Fiber.schedule会抛出RuntimeError: No scheduler is available!对应 cont.c#L2522 的rb_raisetest/fiber/test_scheduler.rb#L8-L14 的test_fiber_without_scheduler验证了这一点。3.2 创建阻塞上下文你也可以显式创建阻塞执行上下文Fiber.new(blocking: true) do # 不会使用调度器 sleep(n) endFiber#blocking?可查询当前 fiber 是否阻塞cont.c#L3010-L3019非阻塞返回false阻塞返回1。此外还有Fiber.blocking { |fiber| ... }用于在块执行期间临时强制当前 fiber 变为阻塞cont.c#L2985-L3008——若当前 fiber 已是阻塞态则近乎 no-op否则在块执行期间临时提高线程的阻塞计数。官方建议除非你正在实现调度器否则应尽量避免创建阻塞上下文Fiber.blocking则常用于在调度器内部把少量必须阻塞的操作如底层系统调用包起来。四、IO 与调度器的交互4.1 默认非阻塞的 I/O默认情况下I/O 是非阻塞的。但并非所有操作系统都支持非阻塞 I/O——Windows 是典型例外其 socket I/O 可以非阻塞但 pipe I/O 是阻塞的。只要满足两个条件——存在调度器且当前线程处于非阻塞状态——IO 操作就会调用调度器走io_wait/io_read/io_write等 hook。4.2 IO#close中断阻塞操作的完整机制IO#close会中断该 IO 上的所有阻塞操作其过程远比关闭文件描述符复杂当线程调用IO#close时先尝试中断所有阻塞在该 IO 上的线程或 fiber关闭线程会一直等待直到所有被阻塞的线程/fiber 都被妥善中断并从该 IO 的阻塞列表中移除每个被中断的线程/fiber 收到一个IOError并干净地退出阻塞操作只有当所有阻塞操作都被中断清理完毕后才会真正关闭文件描述符——这保证了资源清理的正确性避免了竞态条件。对于调度器管理的 fiber中断过程会调用调度器的rb_fiber_scheduler_fiber_interrupt对应Fiber::Scheduler#fiber_interrupthook必需的。调度器以适合其事件循环实现的方式处理中断、通知 fiberfiber 收到IOError后退出阻塞操作。官方文档用序列图精确刻画了这一过程这一机制在 test/fiber/test_io_close.rb 中有完整测试佐证test_io_close_across_fibers第 19-45 行让一个 fiber 在i.read上阻塞、另一个 fiber 调用i.close最终捕获到IOError且错误消息匹配/closed/test_io_close_blocking_fiber第 78-106 行则验证了外部线程直接关闭 IO 也能以同样方式中断调度器中的阻塞 fiber。值得注意的是测试对mswin|mingw平台做了跳过处理Interrupting a io_wait read is not supported!印证了文档中Windows 下非阻塞 I/O 支持不完整的说明。五、非阻塞上下文中的同步原语官方文档明确以下同步机制在非阻塞上下文有调度器中均可使用且都是fiber 级fiber-specific的MutexMutex类可在非阻塞上下文中使用且与 fiber 绑定ConditionVariable可在非阻塞上下文中使用fiber 级Queue / SizedQueue可在非阻塞上下文中使用fiber 级Thread#join可在非阻塞上下文中使用fiber 级。test/fiber/test_mutex.rb 的test_condition_variable展示了典型用法两个Fiber.schedule的 fiber 通过Thread::MutexThread::ConditionVariable协作一个condition.wait(mutex)等待、一个condition.signal唤醒signalled计数最终为 3——证明这些原本面向线程的原语在 fiber 调度场景下依然成立。test/fiber/test_queue.rb 则验证了Queue#pop带超时与返回值的行为test/fiber/test_thread.rb 的test_thread_join验证了在Fiber.schedule内Thread.new{:done}.value隐式 join能正确返回test_thread_join_timeout第 23-43 行验证了Thread#join(0.1)带超时的 join 不会阻塞事件循环。test/fiber/test_process.rb 则覆盖了process_waithook 对应的Process.wait、system与fork场景。六、参考实现一个可运行的最小调度器仓库 test/fiber/scheduler.rb 提供了一个为测试目的编写、刻意简化的完整调度器实现是理解 hook 如何落地的绝佳教材。其头部注释坦诚地说明了局限用IO.select实现、对同一 fd 的多次重叠wait处理不完善生产级调度器应基于 epoll/kqueue例如io-eventgem。但这不妨碍我们从中提炼出调度器的核心骨架。6.1 事件循环run 与 run_once调度器的心脏是事件循环。run在Thread.handle_interrupt保护下循环执行run_once直到所有读、写、等待、阻塞集合为空def run Thread.handle_interrupt(::SignalException :never) do while readable.any? or writable.any? or waiting.any? or blocking.any? run_once break if Thread.pending_interrupt? end end endrun_once用IO.select(readable.keys [urgent.first], writable.keys, [], next_timeout)同时等待 IO 就绪、超时与紧急唤醒管道并把就绪的 fiber 通过fiber.transfer(events)恢复selected.each do |fiber, events| fiber.transfer(events) end6.2 block / unblock最基础的协作契约block由Thread::Mutex#lock、Thread::Queue#pop、Thread::SizedQueue#push等阻塞操作触发让当前 fiber 让出控制权带超时时把 fiber 记入waiting否则记入blockingdef block(blocker, timeout nil) fiber Fiber.current if timeout waiting[fiber] current_time timeout begin fiber.transfer ensure waiting.delete(fiber) # 防止 unblock 先于超时到达导致残留 end else blocking[fiber] true begin fiber.transfer ensure blocking.delete(fiber) end end endunblock则由其他线程或 fiber 调用要求线程安全。参考实现用一个互斥锁把待唤醒 fiber 放入ready然后向urgent管道写一个字节——事件循环正阻塞在IO.select上这一写会立刻唤醒它def unblock(blocker, fiber) lock.synchronize do ready fiber end io urgent.last io.write_nonblock(.) endkernel_sleep直接委托给block(:sleep, duration)io_wait则按IO::READABLE/IO::WRITABLE位掩码把当前 fiber 注册进readable/writable后再让出。fiberhook 则创建blocking: false的 fiber 并立即transfer启动实现Fiber.schedule的立即执行语义def fiber(block) fiber Fiber.new(blocking: false, block) fiber.transfer return fiber end6.3 进阶 hook 的参考实现同一文件还给出了其他 hook 的务实解法process_waitProcess.wait、system、反引号命令触发起一个线程执行Process::Status.wait(pid, flags)并取回结果让阻塞等待不卡住事件循环address_resolveAddrinfo.getaddrinfo触发同样用线程包装因为 libc 的getaddrinfo是阻塞的io_select起线程执行IO.selecttimeout_afterTimeout.timeout触发创建一个休眠duration的 fiber到期后向目标 fiberraise(klass, message)fiber_interrupt把FiberInterrupt内部包装了fiber.raise(exception)压入ready并写管道唤醒事件循环scheduler_close / close调度器离开作用域时先self.run跑完所有剩余任务再关闭紧急管道并freeze防止误改。测试文件底部还定义了若干故意写坏的调度器子类用于验证健壮性BrokenUnblockSchedulerunblock抛异常、SleepingUnblockSchedulerunblock中睡眠改变线程状态、SleepingBlockingSchedulerkernel_sleep中先阻塞睡眠——它们被 test/fiber/test_scheduler.rb 等测试用来确认 Ruby 在调度器行为异常时不会死锁或挂死。七、在仓库中验证与深入学习你可以直接在本仓库中运行这些测试来观察调度器行为例如在仓库根目录执行ruby -Itest/fiber test/fiber/test_scheduler.rb ruby -Itest/fiber test/fiber/test_io_close.rb ruby -Itest/fiber test/fiber/test_mutex.rb ruby -Itest/fiber test/fiber/test_thread.rb ruby -Itest/fiber test/fiber/test_process.rb进一步深入源码时推荐按以下线索阅读cont.cFiber 的完整 C 实现重点看Fiber.new/resume/yield的文档注释#L3147-L3411附近与Fiber.schedule、Fiber.set_scheduler、Fiber.blocking、Fiber.blocking?的实现test/fiber/scheduler.rb可直接复制改造的最小调度器test/fiber 目录test_scheduler.rb、test_io_close.rb、test_mutex.rb、test_queue.rb、test_thread.rb、test_process.rb、test_timeout.rb、test_sleep.rb分别覆盖各 hook 与各同步原语在非阻塞上下文中的行为是编写自定义调度器时最好的行为规范。结语Fiber 与 Fiber::Scheduler 共同构成了 Ruby 面向高并发 I/O 的基础设施Fiber 提供轻量的协作式上下文切换Scheduler 提供可插拔的阻塞拦截层。理解yield/resume/transfer的切换语义、14 个 hook 的职责边界、block/unblock的协作契约以及IO#close的中断时序是写出可靠自定义调度器的前提。仓库中这份官方文档配合cont.c源码与完整测试套件是学习这一机制的权威起点——对照 doc/language/fiber.md 逐项实现 hook再以 test/fiber/scheduler.rb 为参照跑通测试你就能从会用 Fiber进阶到写出自己的事件循环。【免费下载链接】rubyThe Ruby Programming Language项目地址: https://gitcode.com/GitHub_Trending/ru/ruby创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考