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 与 transfer
Fiber 执行用户提供的代码块。在块执行期间,可以调用Fiber.yield或Fiber.transfer切换到其他 Fiber;Fiber#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-L2603):
Fiber.scheduler:返回当前线程最后设置的调度器,未设置时为nil;Fiber.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注意:未来可能引入更多 hook,Ruby 将采用**特性检测(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_required(test/fiber/test_scheduler.rb#L109-L121)进一步验证:即使实现了block/unblock/io_wait/kernel_sleep,若缺少fiber_interrupt,Fiber.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.schedule:
Fiber.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_raise,test/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(必需的)。调度器以适合其事件循环实现的方式处理中断、通知 fiber,fiber 收到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)的:
- Mutex:
Mutex类可在非阻塞上下文中使用,且与 fiber 绑定; - ConditionVariable:可在非阻塞上下文中使用,fiber 级;
- Queue / SizedQueue:可在非阻塞上下文中使用,fiber 级;
- Thread#join:可在非阻塞上下文中使用,fiber 级。
test/fiber/test_mutex.rb 的test_condition_variable展示了典型用法:两个Fiber.schedule的 fiber 通过Thread::Mutex+Thread::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,否则记入@blocking:
def 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_wait(
Process.wait、system、反引号命令触发):起一个线程执行Process::Status.wait(pid, flags)并取回结果,让阻塞等待不卡住事件循环; - address_resolve(
Addrinfo.getaddrinfo触发):同样用线程包装,因为 libc 的getaddrinfo是阻塞的; - io_select:起线程执行
IO.select; - timeout_after(
Timeout.timeout触发):创建一个休眠duration的 fiber,到期后向目标 fiberraise(klass, message); - fiber_interrupt:把
FiberInterrupt(内部包装了fiber.raise(exception))压入@ready并写管道唤醒事件循环; - scheduler_close / close:调度器离开作用域时先
self.run跑完所有剩余任务,再关闭紧急管道并freeze防止误改。
测试文件底部还定义了若干故意写坏的调度器子类用于验证健壮性:BrokenUnblockScheduler(unblock抛异常)、SleepingUnblockScheduler(unblock中睡眠,改变线程状态)、SleepingBlockingScheduler(kernel_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.c:Fiber 的完整 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),仅供参考