
Ray 内核 RPC 容错机制详解幂等性设计、Retryable gRPC Client 与故障注入测试【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray导读Ray 是一个 AI 计算引擎其核心分布式运行时Ray Core内部大量依赖 gRPC 完成进程间通信Core Worker 与 Raylet 之间、Raylet 与 GCSGlobal Control Service之间几乎每一次调度决策都建立在一系列 RPC 之上。在真实生产环境中网络瞬时故障transient network error随时可能发生因此 Ray 要求所有新增的 Core RPC 都必须具备容错能力统一使用 retryable gRPC client并在设计上优先保证幂等性idempotency。本文以RequestWorkerLease请求 Worker 租约为完整案例讲解 Ray 如何通过引入LeaseID实现租约请求去重、如何发现并修复长轮询long-pollingRPC 的隐藏幂等性缺陷随后深入剖析 retryable gRPC client 的重试队列、通道状态检查、指数退避与 GCS 节点状态检查机制最后介绍 Ray Core 用于验证 RPC 容错性的三层测试体系C 单元测试、Python 集成测试、混沌网络发布测试。读完本文你将掌握在 Ray 中设计容错 RPC 的关键原则并能用RAY_testing_rpc_failure等工具验证自己的实现。一、容错与幂等Ray Core 所有 RPC 的设计底线Ray 官方文档rpc-fault-tolerance.rst给出了明确的工程规范所有加入 Ray Core 的 RPC 都应当是容错的并使用 retryable gRPC client。理想情况下这些 RPC 应该是幂等的即便无法做到幂等也必须把不幂等这一事实文档化并让客户端在重试时把这一点考虑进去。什么是幂等考虑一个向文件写入 hello 的函数首次调用写入 hello重试时又写入一次 hello最终文件内容变成 hellohello——这就是不幂等的。要让其幂等可以在再次写入前先检查文件内容确保多次完全相同调用后的可观察状态与单次调用后的状态一致。对于 RPC 而言网络层会做透明重试这是 retryable client 的核心行为但网络层并不知道上层业务逻辑是否允许重复执行。因此幂等性必须由 RPC 的服务端处理器HandleX自己保证当客户端因为请求或响应丢失而重试时服务端要能识别出这是一次重试并返回与首次执行等价的结果而不是重复执行副作用。从 node_manager.proto 可以看到 Ray 的调度 RPC 定义其中RequestWorkerLeaseRequest通过LeaseSpec携带资源描述RequestWorkerLeaseReply返回被租用 Worker 的地址、资源映射以及各种调度失败类型资源不足、placement group 被移除、runtime env 配置失败、任务被取消等——这正是下面案例研究的对象。二、案例研究RequestWorkerLease 的幂等化改造2.1 问题为何它曾经无法重试在修复之前RequestWorkerLease无法做成可重试的因为 Raylet 端的处理器不幂等。其根本原因在于租约的生命周期一旦租约被授予该 Worker 及其资源就被视为已占用状态直到客户端调用ReturnWorker归还。在ReturnWorker到来之前这个 Worker 和它的资源永远不会回到可用 Worker/资源池中。而 Raylet 无法区分原始的租约请求和它的重试请求把两者都当作全新的租约请求来处理。文档给出了一个典型故障序列OwnerCore Worker通过RequestWorkerLease向 Raylet 请求一个新的 Worker 租约响应在传输中丢失Raylet → OwnerOwner 重试RequestWorkerLease请求租约结果原始请求和重试请求各被授予一套资源和 Worker——资源被重复分配。正确行为应该是当重试到达时Raylet 能识别出这是一次重试并把已租出的 Worker 地址转发给 Owner而不是再授予一份新租约。2.2 解决方案引入 LeaseID 做去重解决方案是在RequestWorkerLeaseRequest中引入一个唯一标识符LeaseID由发起方在请求中生成用于对到达的租约请求做去重。LeaseID定义在 common.proto 的LeaseSpec中bytes lease_id 1并随LeaseSpec一起通过RequestWorkerLeaseRequest.lease_spec传递。租约一旦授予就会记录在一个leased_workers映射表中该表将LeaseID映射到具体的 Worker。当一个新的租约请求到达时如果其LeaseID已经存在于leased_workers表中系统就知道这是一次重试直接返回已租用 Worker 的地址而不再重复授予资源。当前仓库中的实现印证了这一点。在 node_manager.cc 的HandleRequestWorkerLease第 1879 行起中处理器首先从请求中解析LeaseID然后检查去重逻辑若leased_workers_已包含该lease_id则记录日志并直接复用已授予的 Worker 信息进程 PID、IP、端口、WorkerID、NodeID构造响应立即回复不再走调度流程node_manager.cc若该lease_id出现在cancelled_lease_tombstones_已取消租约墓碑表中则说明可能是CancelWorkerLease因消息乱序先于本请求到达直接以SCHEDULING_CANCELLED_INTENDED失败类型回复否则才作为全新的租约请求进入调度与资源分配流程。2.3 隐藏问题长轮询Long-pollingRPC网络瞬时错误可能在任意时刻发生。对大多数 RPC 来说它们在一次 I/O 上下文执行内就能完成因此只需防护请求失败或响应失败两种情况即可。但有一小类 RPC 是长轮询式的HandleX函数执行后并不会立即给客户端回包而是要依赖未来某个状态变化来触发响应。RequestWorkerLease正是这种类型租约必须等所有参数args被拉取完成后才能授予因此在参数拉取结束之前系统无法回复客户端。这里就出现了一个棘手场景客户端在 Raylet 侧服务端逻辑还在执行时断开并发送重试。由于leased_workers表只跟踪已授予的租约、不跟踪正在授予中的租约重试请求会通过幂等检查导致 Raylet 无法去重进而在lease_dependency_manager中触发RAY_CHECK崩溃。文档给出的故障序列Owner 通过RequestWorkerLease请求新的 Worker 租约Raylet 正在异步拉取该租约所需的参数argsOwner 重试RequestWorkerLease租约尚未授予重试通过了幂等检查Raylet 去重失败由于 Raylet 尝试为同一个租约再次拉取参数RAY_CHECK被触发。这个崩溃点在当前源码中依然可见LeaseDependencyManager::RequestLeaseDependencies在向queued_lease_requests_插入租约时会执行RAY_CHECK(inserted.second) Lease depedencies can be requested only once per lease.lease_dependency_manager.cc——同一个租约的依赖只能请求一次重复请求就会触发检查失败。最终修复方案是把服务端逻辑可能仍在执行这一点纳入设计跟踪租约在整个授予过程中的各个阶段阶段状态机保证在任何阶段系统都能对请求去重。这对所有长轮询 RPC 都是一个重要警示客户端重试不一定会等待响应发出因此幂等性检查绝不能只覆盖请求已完成这一种状态而必须覆盖请求正在处理中的整个时间窗口。三、Retryable gRPC Client 工作原理Ray 的 RPC 容错项目更新了 retryable gRPC client所有核心客户端Core Worker client、Raylet client、GCS client的 RPC 都通过它发送。其核心头文件注释见 retryable_grpc_client.h类主体在 retryable_grpc_client.cc。3.1 工作流程概览retryable gRPC client 的工作方式如下RPC 通过 retryable gRPC client 发送CallMethod模板retryable_grpc_client.h会把每次调用包装为一个RetryableGrpcRequest增加活跃请求计数后立即发出。遇到 gRPC 瞬时网络错误时入队请求回调返回失败后会先判断该失败是否为可重试状态。可重试状态的定义在 grpc_util.hIsGrpcRetryableStatus仅对UNAVAILABLE和UNKNOWN两种错误码返回 true这两种状态通常表示对端可能已宕机。满足条件则把回调推进重试队列pending_requests_否则直接回调给上层。周期性执行若干检查决定何时重发队列中的请求。3.2 三类周期性检查1廉价的 gRPC 通道状态检查调用channel_-GetState(false)获取 gRPC 通道连接状态retryable_grpc_client.cc判断系统是否可以重新开始发送消息。该检查默认每秒执行一次间隔可通过配置项check_channel_status_interval_milliseconds调整ray_config_def.h 中定义。从CheckChannelStatus的 switch 分支可以看到完整的状态机GRPC_CHANNEL_TRANSIENT_FAILURE/GRPC_CHANNEL_CONNECTING通道仍不可用若超过退避期限则调用server_unavailable_timeout_callback_并递增退避次数、重设下一次超时时间点retryable_grpc_client.ccGRPC_CHANNEL_READY/GRPC_CHANNEL_IDLE通道恢复清空不可用超时状态把队列中的所有请求全部重发并将退避尝试次数归零retryable_grpc_client.ccGRPC_CHANNEL_SHUTDOWN及其他状态直接RAY_LOG(FATAL)理论上不应出现。2可能昂贵的 GCS 节点状态检查如果指数退避期限已过而通道仍然不通系统会调用server_unavailable_timeout_callback_该回调在 retryable_grpc_client.h 声明、在 retryable_grpc_client.cc 触发。这个回调在 client pool 类中设置raylet_client_poolraylet_client_pool.cccore_worker_client_poolcore_worker_client_pool.cc。回调逻辑是先检查该 client 是否订阅了节点状态更新若已订阅则查询本地 subscriber 缓存中是否收到了来自 GCS 的节点死亡通知若 client 未订阅或缓存中没有该节点的状态则直接向 GCS 发起一次 RPC 查询。对 GCS client 而言server_unavailable_timeout_callback_一旦被调用就直接杀死进程见 rpc_client.h 附近实现这发生在gcs_rpc_server_reconnect_timeout_s秒默认 60之后ray_config_def.h。这也解释了文档中的说明GCS client 没有最大退避上限因为它的兜底策略是超时就终止进程。3每个 RPC 的超时检查CheckChannelStatus在每次检查开始时会先清理所有已超时的挂起请求以Status::TimedOut失败它们retryable_grpc_client.cc。该超时是每个 RPC 可定制的但在 Ray Core 的实践中它被功能性地禁用了Core Worker client 对每个 RPC 的超时参数都固定传-1表示无限见 core_worker_client.h。因此实际语义是RPC 可以无限期排队重试直到通道恢复、节点死亡确认或进程被杀。3.3 指数退避与队列重发退避递增每增加一次失败 RPC指数退避周期就会增大retryable_grpc_client.cc 中通过ExponentialBackoff::GetBackoffMs计算详见 exponential_backoff.h。退避与失败 RPC 的类型无关对所有请求一视同仁。退避上限可通过以下配置自定义core_worker_rpc_server_reconnect_timeout_max_s默认 60 秒ray_config_def.h作用于 Core Worker clientraylet_rpc_server_reconnect_timeout_max_s默认 60 秒ray_config_def.h作用于 Raylet client。如前述GCS client 没有最大退避——它以进程终止作为最终兜底。退避重置一旦通道检查成功指数退避周期即被重置队列中所有 RPC 全部重发attempt_number_ 0见 retryable_grpc_client.cc。节点死亡确认后的行为如果系统通过订阅或直接查询 GCS 成功收到了节点死亡通知就会销毁对应的 RPC client销毁过程会把每个待处理回调投递到 I/O 上下文并携带gRPC Disconnected 错误Status::Disconnected(GRPC client is shut down.)见 retryable_grpc_client.cc。上层应用代码必须把Disconnected当作对端节点已死、重试无意义的信号来处理。四、使用 retryable client 必须注意的三个要点4.1 队列是按客户端而非按 RPC 类型的每个 retryable gRPC client 唯一对应一个对端 client 实例——Core Worker client 以WorkerID区分Raylet client 以NodeID区分——而不是对应某种 RPC 类型。假设你先提交 RPC A它因瞬时网络错误失败入队随后向同一个 client 提交 RPC B同样失败入队。那么队列里会有两个元素先 A 后 B。没有按 RPC 划分的独立队列只有按 client 划分的统一队列。4.2 超时是客户端级的会累加队列中的每个超时必须等待前一个超时完成。如果 RPC A 和 RPC B 在很短时间内先后提交那么 A 总共等待 1 秒而 B 要等待 1 2 3 秒假设退避依次为 1 秒、2 秒。不同 RPC 之间没有差别都被同等对待。其设计推理是瞬时网络错误不是 RPC 特有的——如果 RPC A 遭遇网络故障可以合理假设发往同一 client 的 RPC B 也会遭遇同样的故障。因此一个 RPC 的等待时间等于队列中所有此前 RPC 的超时之和再加上它自己的超时。4.3 析构函数的行为约束在RetryableGrpcClient的析构函数中系统会把所有挂起 RPC 失败掉通过向它们的 I/O 上下文投递回调实现见 retryable_grpc_client.cc。这些回调理想情况下绝不应修改 client 类如RayletClient持有的状态如果确实必须修改则必须通过某种方式例如弱指针检查 client 是否仍然存活。此外应用代码也必须显式处理Disconnected错误——这正是上面提到的节点已死语义。五、三层测试体系如何验证 RPC 容错性Ray Core 对 RPC 容错与幂等性的验证有三层测试层层递进、互为补充。5.1 第一层C 单元测试对每个 RPC都应有某种形式的 C 幂等性测试调用两次HandleX服务端函数检查两次输出是否一致并且要考虑两次HandleX调用之间状态变化的差异。例如对RequestWorkerLease就专门编写了 C 单元测试来模拟初始租约请求卡在参数拉取阶段时重试到达的场景——正是 2.3 节描述的隐藏问题。相关测试可以在 cluster_lease_manager_test.cc 与 local_lease_manager_test.cc 中找到。5.2 第二层Python 集成测试如果某个 RPC 用 Python API 可以直观测试理想情况下应为它编写一个 Python 集成测试。但有些 RPC 很难用 Python API 完全确定性地测试此时充分的 C 单元测试可以充当很好的替代proxy。因此 Python 集成测试更多是锦上添花nice-to-have同时它们也向用户展示了如何在实际使用中遇到幂等性问题。主要的测试机制是RAY_testing_rpc_failure配置项在 ray_config_def.h 中有完整注释它允许在运行时注入三类故障请求失败request failure不实际发送 RPC立即以 gRPC 错误触发 RPC 回调响应失败response failure正常发送 RPC等服务器响应到达后再以 gRPC 错误触发回调在途失败in-flight failure立即以 gRPC 错误触发回调但同时把 RPC 发送给服务器——对长轮询 RPC 来说重试理想情况下应当在服务器执行服务端代码期间命中服务器从而测试幂等性窗口。配置格式为 JSON 字符串通过环境变量注入。例如为所有 RPC 注入无限次故障并设置 25% 请求失败、50% 响应失败、10% 在途失败export RAY_testing_rpc_failure{*:{num_failures:-1,req_failure_prob:25,resp_failure_prob:50,in_flight_failure_prob:10}}其中*通配符表示对所有 RPC 方法生效设置通配符会覆盖其他方法的具体配置num_failures设为-1表示无限次失败。还可以通过可选的第五、六、七个参数num_lower_bound_req_failures、num_lower_bound_resp_failures、num_lower_bound_in_flight_failures指定前 X 次请求失败、随后 Y 次响应失败、再随 Z 次在途失败保证至少发生 N 次失败之后才回退到概率注入export RAY_testing_rpc_failure{*:{num_failures:-1,req_failure_prob:25,resp_failure_prob:50,in_flight_failure_prob:10,num_lower_bound_req_failures:2,num_lower_bound_resp_failures:3,num_lower_bound_in_flight_failures:1}}此外还有配套开关testing_rpc_failure_avoid_intra_node_failuresray_config_def.h启用后如果服务器和客户端在同一地址同节点则不注入故障便于把故障注入限定在跨节点通信上。5.3 第三层混沌网络发布测试Chaos Network Release TestsIP 表黑屏IP table blackout方案通过 SSH 进入每个节点短暂5 秒屏蔽其 IP 表来模拟瞬时网络错误。IP 表脚本在后台运行在测试脚本执行期间周期性每 60 秒制造一次网络黑屏。相比真实拔线/断网这种方案对集群的扰动是可控、可恢复的。文档特别说明最初曾考虑过 Amazon FISFault Injection Service但 FIS 有 60 秒的最小执行时间这个时长会触发节点因配置如心跳超时而被判定死亡且问题难以调试相比之下IP 表方案更简单、更灵活因此被最终采用。在 Ray 的 core release 测试中这一方案已被应用到所有已有的混沌发布测试里。六、给 RPC 开发者的检查清单综合文档与源码当你在 Ray Core 中新增一个 RPC 时建议按以下清单自查统一走 retryable client所有 RPC 必须通过RetryableGrpcClient发送而不是裸 gRPC stub新增 client 方法时可借助INVOKE_RETRYABLE_RPC_CALL/VOID_RETRYABLE_RPC_CLIENT_METHOD宏retryable_grpc_client.h快速接入。设计幂等键为请求设计唯一标识如RequestWorkerLease的LeaseID服务端以标识 → 处理结果的映射表完成去重。覆盖处理中窗口如果你的 RPC 是长轮询式的去重逻辑必须覆盖服务端逻辑仍在执行的整个阶段而非仅覆盖已完成状态否则客户端重试会穿透幂等检查。明确失败语义区分可重试错误UNAVAILABLE/UNKNOWN与不可重试错误对Disconnected错误要有明确的处理路径节点已死重试无意义。补全三层测试至少编写调用两次HandleX的 C 幂等单元测试能用 Python 触发的场景补集成测试借助RAY_testing_rpc_failure注入三类故障最后通过 IP table blackout 混沌测试验证真实网络环境下的行为。文档化非幂等性如果某个 RPC 确实无法幂等必须在代码和文档中明确说明并确保客户端在重试时已考虑该限制。遵循这些原则新增的 RPC 才能在节点崩溃、网络抖动、消息重排等真实故障下不破坏 Ray 的调度正确性与资源一致性。参考路径索引官方文档rpc-fault-tolerance.rstretryable client 头文件与注释retryable_grpc_client.hretryable client 实现retryable_grpc_client.cc可重试错误判定grpc_util.hRequestWorkerLease协议定义node_manager.protoLeaseID定义common.protoRaylet 幂等去重实现node_manager.cc依赖管理器幂等检查lease_dependency_manager.cc配置项超时、退避上限、故障注入ray_config_def.h相关单元测试cluster_lease_manager_test.cc、local_lease_manager_test.cc【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考