尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Ray 分布式 multiprocessing.Pool:用一行 import 将 Python 多进程程序扩展到集群

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),仅供参考
返回列表