Ray 分布式 multiprocessing.Pool:用一行 import 将 Python 多进程程序扩展到集群
Ray 分布式 multiprocessing.Pool用一行 import 将 Python 多进程程序扩展到集群【免费下载链接】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 提供的ray.util.multiprocessing.Pool一个与 Python 标准库multiprocessing.Pool保持 API 兼容的分布式线程池替代品。它用 Ray Actor 取代本地进程来执行任务让你可以把原本运行在单机上的multiprocessing.Pool程序以几乎零改动的代价扩展到多节点 Ray 集群。读完本文你将掌握它的快速上手方式、完整构造参数、连接集群的三种途径、底层基于 Actor 的实现原理以及close/terminate/join与异步结果的管理技巧。概述为什么需要分布式 multiprocessing.PoolPython 标准库multiprocessing.Pool通过本地进程池并行执行任务其能力上限受限于单机 CPU 数量与内存。Ray 对该 API 的移植思路非常直接不再为每个 worker 启动一个本地进程而是为每个 worker 创建一个 Ray Actor。这些 Actor 既可以是本机进程也可以分布在集群的多个节点上从而让池化并行从单机平滑扩展到集群。从源码看这一设计体现在 python/ray/util/multiprocessing/pool.py 中的Pool类其 docstring 明确写道A pool of actor processes that is used to process tasks in parallel用于并行处理任务的 Actor 进程池。Ray 的 Actor 本身运行在独立进程中因此每个任务仍然享有进程级隔离与并行度但调度与资源管理全部交给 Ray 运行时统一完成。快速开始一行替换 multiprocessing.Pool首先安装 Ray参见 安装指南可执行pip install -U ray[default]然后在代码中把multiprocessing.Pool换成ray.util.multiprocessing.Poolfrom ray.util.multiprocessing import Pool def f(index): return index pool Pool() for result in pool.map(f, range(100)): print(result)运行这段代码时第一次创建Pool会自动启动一个本地 Ray 集群并把 100 个任务分发到该集群的 Actor 上执行。也就是说即使在单机上你也无需显式调用ray.init()——Pool构造器会在必要时自动完成 Ray 的初始化。官方文档指出multiprocessing.Pool的完整 API 目前均受支持apply、map、starmap、imap、imap_unordered以及各自的异步变体等。这意味着你现有使用multiprocessing.Pool的代码往往只需要修改 import 一行即可获得分布式能力。注意Pool构造器中的context参数在 Ray 实现中被忽略。Ray 完全接管进程的初始化方式传入非None值时会记录一条警告日志。这一行为可以在 pool.py 的__init__中看到源码对context参数调用log_once(context_argument_warning)并打印 The context argument is not supported using ray. Please refer to the documentation for how to control ray initialization.。构造参数详解ray.util.multiprocessing.Pool的完整构造函数签名如下见 pool.pyPool( processesNone, initializerNone, initargsNone, maxtasksperchildNone, contextNone, ray_addressNone, ray_remote_argsNone, )各参数含义参数默认值说明processesNone池中 Actor 进程的数量。若已存在运行中的 Ray 集群默认取集群的 CPU 核数否则取本机 CPU 核数。传入的值必须大于 0且不能超过集群可用 CPU 数否则抛出ValueErrorinitializerNone每个 Actor 启动时执行一次的初始化函数对应测试见 test_multiprocessing.py 的 test_initializerinitargsNone传递给initializer的位置参数元组maxtasksperchildNone每个 Actor 最多执行的任务数达到后该 Actor 会被终止并替换为新 Actor对应测试见 test_multiprocessing_standalone.py 的 test_maxtasksperchildcontextNone仅用于 API 兼容Ray 实现中忽略并给出警告ray_addressNone要连接的 Ray 集群地址None时在本机启动新本地集群ray_remote_argsNone配置构成池的 Ray Actor 的参数如num_gpus、resources、max_retries等最终透传给ray.remote的 options其中processes的校验逻辑位于 pool.py 的_init_ray若processes 0抛出ValueError(Processes in the pool must be 0.)若processes大于集群 CPU 数则抛出ValueError提示集群 CPU 不足。测试 test_multiprocessing_standalone.py 专门验证了在只有 4 个 CPU 的集群上创建 8 进程的池会报错这一行为。maxtasksperchild的实现值得注意在_run_batch中每个 Actor 维护一个已执行任务计数当计数达到上限时会先调用__ray_terminate__.remote()优雅停止旧 Actor等待其排空已提交任务再创建新的 Actor 补位从而实现 worker 定期轮换。连接 Ray 集群的三种方式原文档指出要让Pool连接到一个运行中的 Ray 集群有两条等价途径设置RAY_ADDRESS环境变量或向Pool构造函数传入ray_address关键字参数。此外还可以先手动调用ray.init()再创建Pool。综合起来一共有三种方式from ray.util.multiprocessing import Pool # 方式一不指定任何地址——启动一个新的本地 Ray 集群。 pool Pool() # 方式二连接运行中的 Ray 集群当前节点作为 head 节点。 # 等价于设置环境变量 RAY_ADDRESSauto。 pool Pool(ray_addressauto) # 方式三连接运行中的 Ray 集群head 节点为远程节点。 # 等价于设置环境变量 RAY_ADDRESSip_address:port。 pool Pool(ray_addressip_address:port)你还可以先手动启动 Ray 再创建Pool这样能够利用ray.init()支持的全部配置选项import ray ray.init(addressauto, num_cpus8) # 或任何 ray.init() 支持的配置 pool Pool()底层连接优先级Pool的 Ray 初始化逻辑_init_ray遵循如下优先级见 pool.py若 Ray 已初始化ray.is_initialized()为真直接复用现有运行时不做任何初始化否则ray_address参数优先于RAY_ADDRESS环境变量若ray_address为None但环境中存在RAY_ADDRESS或检测到默认的 Ray 地址ray._private.utils.read_ray_address()返回值非空则按集群模式初始化以上条件都不满足时回退到本地模式调用ray.init(num_cpusprocesses)启动本地集群。还有一个特殊取值值得注意ray_addresslocal或RAY_ADDRESSlocal会强制启动一个新的本地 Ray 集群并将集群 CPU 数设置为processes。这一行为在 test_multiprocessing_standalone.py 的 test_connect_to_ray 中有专门验证。关于如何启动和管理多节点 Ray 集群可以参考仓库中的 集群关键概念、集群快速上手 以及 ray start CLI 说明。底层原理Actor 池、分批调度与结果收集为了让文章不止停留在能用这里结合源码剖析Pool的内部工作方式。每个 worker 就是一个 Actor构成池的 worker 是 PoolActor声明为ray.remote(num_cpus0)ray.remote(num_cpus0) class PoolActor: def __init__(self, initializerNone, initargsNone): if initializer: initargs initargs or () initializer(*initargs) def ping(self): # 用于等待该 Actor 初始化完成。 pass def run_batch(self, func, batch): results [] for args, kwargs in batch: ... try: results.append(func(*args, **kwargs)) except Exception as e: results.append(PoolTaskError(e)) return results要点Actor 申请0 个 CPU避免与任务的资源配额冲突这也是它能与ray_remote_args中自定义资源请求共存的原因Pool.__init__阶段会启动processes个 Actor 并逐个调用ping.remote()等待其就绪见_start_actor_pool任务执行失败时不会直接抛异常中断 Actor而是把异常包装成PoolTaskError放入结果列表由调用侧统一还原——这正是pool.map能像标准库一样把任务异常传播给调用者的关键。轮询分发与分块chunking任务通过 round-robin轮询方式分发给各个 Actor见_next_actor_index。对于map系列输入 iterable 会被切分成块chunk每块作为一个批任务交给某个 Actor 的run_batch执行分块大小由_calculate_chunksize决定def _calculate_chunksize(self, iterable): chunksize, extra divmod(len(iterable), len(self._actor_pool) * 4) if extra: chunksize 1 return chunksize即ceil(len(iterable) / (4 × Actor 数))与标准库的分块策略一致可以在调度开销和负载均衡粒度之间取得平衡。测试 test_map 验证了 100 个任务在 4 个 worker 间较为均匀地分布每个 PID 处理超过 20 个任务。惰性迭代imap / imap_unordered与map一次性提交全部任务不同imap系列采用惰性提交策略IMapIterator初始化时只为每个 Actor 提交一批任务之后每消费一个结果、或通过ResultThread收到新批次就绪信号才提交下一批见 IMapIterator。这在输入 iterable 非常大、或单个任务参数占用大量内存时尤其有用。imap返回 OrderedIMapIterator结果按输入顺序返回即使任务完成顺序不同imap_unordered返回 UnorderedIMapIterator结果按完成顺序返回吞吐优先。两个迭代器的next(timeoutNone)均支持超时测试 test_imap_timeout 和 test_imap_unordered_timeout 验证了超时与乱序返回语义。注意imap的 iterable 必须可迭代传入非 iterable 会抛出TypeError见 test_imap_fail_on_non_iterable。结果收集ResultThread 与 AsyncResult异步结果由 ResultThread 这个守护线程负责收集。它维护待就绪的 ObjectRef 列表用ray.wait(unready, num_returns1, timeout0.1)以 100ms 为周期轮询就绪状态源码注释说明该超时是经验值权衡了轮询开销与首个结果的尾延迟。当收到END_SENTINEL哨兵时停止等待随后统一触发成功回调或错误回调。AsyncResult 是暴露给用户的异步接口提供wait(timeoutNone)等待完成超时不抛异常get(timeoutNone)阻塞获取结果超时抛出multiprocessing.TimeoutError该异常直接从multiprocessing导入并在 python/ray/util/multiprocessing/init.py 中导出ready()结果是否就绪successful()所有任务是否成功仅在就绪后可调用否则抛ValueError。回调语义与标准库一致callback仅在所有结果都成功时被调用一次一旦出现首个失败结果则调用error_callback仅一次并跳过常规回调。测试 test_callbacks 还专门把 Ray 与原生multiprocessing.Pool的回调行为做了对比验证。完整任务提交 APIPool完整支持以下任务提交方法与multiprocessing.Pool一一对应方法语义返回apply(func, argsNone, kwargsNone)在任意一个 Actor 上执行一次调用同步返回结果apply_async(func, args, kwargs, callback, error_callback)上述调用的异步版本AsyncResultmap(func, iterable, chunksizeNone)对 iterable 每个元素执行func同步返回结果列表listmap_async(func, iterable, chunksize, callback, error_callback)上述调用的异步版本AsyncResultstarmap(func, iterable, chunksizeNone)同map但解包参数func(*args)liststarmap_async(func, iterable, callback, error_callback)上述调用的异步版本AsyncResultimap(func, iterable, chunksize1)惰性提交按输入顺序产出结果OrderedIMapIteratorimap_unordered(func, iterable, chunksize1)惰性提交按完成顺序产出结果UnorderedIMapIterator示例starmap 与异步回调from ray.util.multiprocessing import Pool def add(x, y): return x y pool Pool(processes4) # starmap解包每个元素作为独立参数 print(pool.starmap(add, [(1, 2), (3, 4)])) # [3, 7] # map_async异步提交 成功回调 result pool.map_async(add, [(1, 2), (3, 4)] if False else range(10)) # 注意 map 会把每个元素作为唯一位置参数传给 func参数解包请用 starmap pool.close() pool.join()生命周期管理close / terminate / join 与上下文管理器Pool的生命周期管理与标准库保持一致三者组合使用close()禁止提交新任务但允许已提交的未完成任务继续执行随后优雅停止各 Actorterminate()立即停止所有未完成任务并杀掉 Actor源码中通过ray.kill(actor)实现见 pool.pyjoin()等待已关闭池中的 Actor 全部退出若池尚未关闭close/terminate均未调用抛出ValueError(Pool is still running)。测试 test_close 与 test_terminate 精确刻画了两者的差异close()后阻塞中的任务仍能完成而terminate()后join()会立即返回、任务结果标记为失败。Pool还实现了上下文管理器协议__enter__/__exit__见 pool.py退出with块时自动调用terminate()from ray.util.multiprocessing import Pool with Pool(processes4) as pool: results pool.map(f, range(100)) # with 块结束后池被自动终止进阶行为与注意事项Ray 初始化对既有运行时的影响测试 test_ray_init 总结了三条行为准则若 Ray 尚未初始化创建Pool会启动本地 Ray 集群CPU 数等于processes若已有本地集群在运行创建Pool不会改动既有集群配置若既有集群 CPU 不足于processes抛ValueError。递归任务与死锁规避Pool支持在任务内部再创建Pool嵌套使用。测试 test_deadlock_avoidance_in_recursive_tasks 验证了递归嵌套的pool.map能够正常返回不会因 Actor 资源占用导致死锁。大对象经由对象存储传输为了减少重复序列化开销任务参数若大于 100 字节会被自动放入 Ray 对象存储并以ObjectRef形式传给 ActorActor 端通过ray.get取回见ray_put_if_needed与ray_get_if_needed。Pool内部维护 list/dict 两个注册表做对象去重_registry/_registry_hashableclose()时统一清理。与 joblib 的协作若环境中安装了 joblibPool会把 joblib 的BatchedCalls转换为 RayBatchedCalls使批量调用中的公共参数只放入对象存储一次节约时间与内存。由于该能力依赖 joblib 的可选导入见 pool.py 顶部未安装 joblib 时相关功能自动退化为原生调用。小结ray.util.multiprocessing.Pool是一个投入产出比极高的分布式能力入口API 完全对齐Python 标准库multiprocessing.Pool迁移成本近乎为零只需改 import本地即刻可用第一次创建Pool自动启动本地 Ray单机体验与标准库一致集群按需扩展通过RAY_ADDRESS环境变量、ray_address参数或先行ray.init()三种方式接入多节点集群实现可追溯Actor 池、分块调度、轮询分发、ResultThread结果收集、异步接口与生命周期管理均可在 python/ray/util/multiprocessing/pool.py 中对应到具体实现配套测试位于 python/ray/tests/test_multiprocessing.py 与 python/ray/tests/test_multiprocessing_standalone.py。如果你的应用目前受限于单机multiprocessing.Pool的规模这是把并行能力平滑扩展到 Ray 集群的最直接路径之一。【免费下载链接】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),仅供参考

相关新闻

Hero 开源库 CHANGELOG 深度解读:从 1.3.0 到 1.6.3 的版本演进与核心源码实现

Hero 开源库 CHANGELOG 深度解读:从 1.3.0 到 1.6.3 的版本演进与核心源码实现

Hero 开源库 CHANGELOG 深度解读:从 1.3.0 到 1.6.3 的版本演进与核心源码实现 【免费下载链接】Hero Elegant transition library for iOS & tvOS 项目地址: https://gitcode.com/gh_mirrors/he/Hero Hero 是面向 iOS 与 tvOS 的优雅转场动画库&#xf…

2026/9/20 23:55:58 阅读更多 →
grok-build v0.2.87 版本解读:自动订阅升级、/docs 导航与按模型推理强度配置

grok-build v0.2.87 版本解读:自动订阅升级、/docs 导航与按模型推理强度配置

grok-build v0.2.87 版本解读:自动订阅升级、/docs 导航与按模型推理强度配置 【免费下载链接】grok-build SpaceXAIs coding agent harness and TUI. Fullscreen, mouse interactive, extensible. 项目地址: https://gitcode.com/gh_mirrors/gr/grok-build …

2026/9/22 1:59:34 阅读更多 →
Atlas 300V 24G部署YOLOv5全流程:昇腾推理加速卡实战指南

Atlas 300V 24G部署YOLOv5全流程:昇腾推理加速卡实战指南

Atlas这个词,做AI模型部署的人这两年应该不陌生。它不是一个单一的硬件,而是华为昇腾整个AI计算平台的代号,从边缘小盒子到服务器推理卡都有覆盖。我手头这块Atlas 300V 24G,就是一张非常典型的PCIe形态AI推理加速卡,常…

2026/9/20 23:54:58 阅读更多 →

最新新闻

手机盖板渲染原理图解:从像素到GPU的最佳实践

手机盖板渲染原理图解:从像素到GPU的最佳实践

手机盖板渲染原理图解:从像素到GPU的最佳实践 看了一堆教程还是不会写项目?这种无力感我太懂了。你盯着屏幕上的精美UI,心里却发慌:这玻璃质感、这光影反射,到底怎么算出来的?别急,今天咱们不整虚的,直接拆解 手机盖板…

2026/9/22 2:02:05 阅读更多 →
微信聊天记录能找回吗:后端架构师视角的数据恢复完整示例

微信聊天记录能找回吗:后端架构师视角的数据恢复完整示例

微信聊天记录能找回吗:后端架构师视角的数据恢复完整示例 刚拿到一个“数据恢复”的面试题,我盯着屏幕上的代码复制了一堆,结果一跑就报 NullPointerException…

2026/9/22 2:02:05 阅读更多 →
5个实战项目揭秘vb编程软件性能瓶颈

5个实战项目揭秘vb编程软件性能瓶颈

5个实战项目揭秘vb编程软件性能瓶颈 VB6 升级 VB.NET 后,API 全变了,老代码跑不动。我在三个电商后台做实战项目时,发现 80% 的卡顿源于数据绑定和循环渲染。 MDN Web Docs 虽主攻 Web,但其 DOM…

2026/9/22 2:02:05 阅读更多 →
3招搞定女孩青春期叛逆源码解析与实战指南

3招搞定女孩青春期叛逆源码解析与实战指南

3招搞定女孩青春期叛逆源码解析与实战指南 看了一堆心理学书籍和育儿教程,还是搞不定家里那位“变脸比翻书还快”的女儿?这种无力感,就像你背熟了Java的API文档,却写不出一个能跑的Spring…

2026/9/22 2:02:05 阅读更多 →
搞懂第四方物流公司技术栈,这份保姆级教程能救命

搞懂第四方物流公司技术栈,这份保姆级教程能救命

搞懂第四方物流公司技术栈,这份保姆级教程能救命 半夜三点,生产环境突然崩了,满屏红色的 StackTrace 像天书一样滚过屏幕。你盯着那些 NullPointerException 和 ConnectionRefused…

2026/9/22 2:02:05 阅读更多 →
搞懂我的世界光影渲染源码,面试不再慌

搞懂我的世界光影渲染源码,面试不再慌

搞懂我的世界光影渲染源码,面试不再慌 面试时被追问光影原理答不上来,那种尴尬谁懂?别怪面试官刁难,是你把《我的世界光影》当成了纯美术资产,没摸透背后的 源码解析 。今天不聊虚的,直接扒开OptiFine和Iris…

2026/9/22 2:01:05 阅读更多 →

日新闻

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天 配置环境就卡半天?别怪机器慢,多半是你没选对工具链。在Java、Go或Python的项目现场, 手写实现…

2026/9/22 0:00:41 阅读更多 →
剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑 面试被问原理答不上来,是不是常态?别慌。很多开发者对着 GitHub 开源仓库里的代码发呆,看似简单实则暗藏玄机。今天这份【剑帝加点】速查手册,直接带你拆解核心实现,把面试必考的原理讲透。…

2026/9/22 0:00:41 阅读更多 →
手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优 复制来的代码跑不通不知道怎么调?别慌,这种“复制粘贴地狱”在开发圈太常见了。尤其是做 图片压缩网站…

2026/9/22 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/21 3:13:20 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/21 2:19:36 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/21 4:51:05 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/21 15:36:51 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/21 15:36:51 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/19 23:35:34 阅读更多 →