资讯动态

Polars 与 Python multiprocessing:为什么必须使用 spawn 而非 fork

发布时间:2026/9/10 23:47:41 来源:尧图企业网站定制
Polars 与 Python multiprocessing为什么必须使用 spawn 而非 fork【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars导读Polars 是一个基于 Rust 编写、默认就启用多线程并发的 DataFrame 查询引擎当你试图在 Python 代码中把multiprocessing与 Polars 一起使用却遭遇进程卡死deadlock或与「multiprocessing 启动方法」相关的报错时根因几乎总是出在 Unix 系统默认的fork启动方式上。本文将说明何时不需要甚至不应该手动引入多进程、何时值得使用multiprocessing并给出可直接套用的spawn推荐写法、完整的复现示例与失败机理分析帮助你安全地把 Polars 放进多进程代码中。结论先行用 spawn不要用 fork如果你把 Python 内置的multiprocessing模块与 Polars 搭配使用时遇到了报错或死锁请先把进程池/进程的启动方法改为spawnfrom multiprocessing import get_context def my_fun(s): print(s) with get_context(spawn).Pool() as pool: pool.map(my_fun, [input1, input2, ...])对应的完整示例源码位于 docs/source/src/python/user-guide/misc/multiprocess.py其中的recommendation片段即上面的推荐写法。之所以刻意使用get_context(spawn)而不是直接依赖全局默认是因为 UnixLinux、BSD 等上的默认启动方法恰恰就是危险的fork。何时不应该使用 multiprocessingPolars 本身已经吃满你的 CPU 核心Polars 从设计之初就是多线程的引擎会把可并行的计算拆到多个线程上执行。文档中给出的两个典型例子是在同一个select语句中请求两个表达式时两者可以并行计算最后才合并结果使用group_by().agg(expr)做分组聚合时每个分组可以被独立并行求值。在上述这些场景下再叠加一层multiprocessing几乎不可能带来额外提速反而会引入进程间通信、序列化与内存复制开销。如果你还在使用 Polars 的 GPU Engine同样应避免手动多进程——两者同时运行时会在系统内存和算力上互相竞争最终反而降低性能。相关的其它优化手段可参考 用户指南-惰性计算与优化。多线程机制在源码中的证据在底层线程数是可配置的。Rust 侧 crates/polars-config/src/lib.rs 定义了环境变量POLARS_MAX_THREADS其默认值来自std::thread::available_parallelism()——也就是默认采用机器可用的全部并行能力max_threads字段见 crates/polars-config/src/lib.rs直接控制了 Polars 用于 CPU 密集型工作的线程池大小。Python 侧的说明见 py-polars/src/polars/meta/thread_pool.py且该变量在仓库的测试中被大量使用例如 py-polars/tests/unit/io/test_lazy_parquet.py 通过POLARS_MAX_THREADS1把并行度降为 1 以测试串行路径。这意味着当你手动创建大量进程时每个进程还会各自再拉起 Polars 的线程池造成典型的线程/进程叠加放大进一步印证了「多数情况下不需要手动多进程」的结论。何时值得使用 multiprocessing尽管 Polars 是多线程的其它库未必如此。当计算瓶颈落在某个单线程的第三方库上而问题本身又天然可并行例如对一批独立的文件逐一执行单线程库的处理时用multiprocessing绕过 GIL、并行调用这些单线程库就是合理的提速手段。这里的判断标准是瓶颈必须是「别的库」而非 Polars 本身。默认配置下 multiprocessing 的问题三种启动方法与问题摘要Python 的 multiprocessing 文档 提供了三种创建进程池的方法spawnforkforkserver其中对fork的描述截至 2022-10-15是父进程通过 os.fork() 复制出 Python 解释器。子进程开始时与父进程几乎完全相同父进程的所有资源都被子进程继承。注意安全地 fork 一个多线程进程是存在问题的。仅 Unix 可用并且是 Unix 上的默认方法。由此可得到本文最核心的推论Polars 是多线程的因此不能与fork组合使用。只要你在 UnixLinux、BSD 等上且没有显式覆盖fork就是默认启动方法。你之所以以前可能没踩过这个坑往往是因为纯 Python 代码及大多数 Python 库基本是单线程的fork 一个单线程进程没有风险或者你运行在 Windows / macOS 上——这些平台根本不提供fork方法macOS 上直到 Python 3.7 还提供之后移除。因此应当改用spawn或forkserver。其中spawn在所有平台都可用、且最安全是官方推荐方法。复现示例fork 导致死锁问题出在fork对父进程地址空间的整体复制上。下面这个稍作修改的示例改编自 Polars issue tracker原文见 docs/source/src/python/user-guide/misc/multiprocess.py 的example1片段import multiprocessing import polars as pl def test_sub_process(df: pl.DataFrame, job_id): df_filtered df.filter(pl.col(a) 0) print(fFiltered (job_id: {job_id}), df_filtered, sep\n) def create_dataset(): return pl.DataFrame({a: [0, 2, 3, 4, 5], b: [0, 4, 5, 56, 4]}) def setup(): # some setup work df create_dataset() df.write_parquet(/tmp/test.parquet) def main(): test_df pl.read_parquet(/tmp/test.parquet) for i in range(0, 5): proc multiprocessing.get_context(spawn).Process( targettest_sub_process, args(test_df, i) ) proc.start() proc.join() print(fExecuted sub process {i}) if __name__ __main__: setup() main()把上面get_context(spawn)替换成get_context(fork)后程序就会死锁。其机理如下fork等价于调用os.fork()即 POSIX 标准 中定义的系统调用创建出的进程只有单个线程。如果多线程进程调用 fork()新进程会包含调用线程的一个副本及其完整地址空间其中可能包括互斥锁及其它资源的状态。因此为避免错误子进程在被 exec 之前只能执行 async-signal-safe 的操作。而spawn会启动一个全新的、干净的 Python 解释器不继承父进程任何锁的状态。回到示例本身pl.read_parquet读取文件时需要对文件加锁随后os.fork()被调用把父进程状态包括已被持有的文件锁一并复制进子进程。于是每个子进程都「继承」了一把处于已获取状态的锁然后永远等待这把永远不会被释放的文件锁——这就是死锁的来源。因此演示代码里虽然用的是spawn也能说明问题但只要换成fork死锁立刻出现。最棘手的部分fork「偶尔」能用这类 bug 真正折磨人的地方在于fork并非必定失败。把上面例子中的pl.read_parquet调用去掉即不先加锁再 fork改为直接在父进程内构造 DataFrameimport multiprocessing import polars as pl def test_sub_process(df: pl.DataFrame, job_id): df_filtered df.filter(pl.col(a) 0) print(fFiltered (job_id: {job_id}), df_filtered, sep\n) def create_dataset(): return pl.DataFrame({a: [0, 2, 3, 4, 5], b: [0, 4, 5, 56, 4]}) def main(): test_df create_dataset() for i in range(0, 5): proc multiprocessing.get_context(fork).Process( targettest_sub_process, args(test_df, i) ) proc.start() proc.join() print(fExecuted sub process {i}) if __name__ __main__: main()这个版本见源码example2片段可以正常工作。正因如此在大型代码库中排查这类问题会非常痛苦某个看似无关的改动例如让多进程调用之前的某个 I/O 操作多了一把锁就可能瞬间打破原本「碰巧能跑」的多进程代码。因此总的原则是除非存在无法以其它方式满足的特殊需求否则永远不要在多线程库如 Polars上使用fork启动方法。fork 的优缺点为什么它还是默认方法看过上面的示例你可能会问既然fork这么危险为什么 Python 还把它保留为 Unix 默认方法文档给出了四点原因历史原因spawn直到 Python 3.4 才被加入而fork从 Python 2.x 时代就存在。作为长期默认值改变它会破坏大量现有代码的兼容性。spawn/forkserver存在限制它们要求传入的参数必须可 pickle 序列化fork则没有这个约束详见 Python multiprocessing 文档中关于 spawn 与 forkserver 的说明。在需要把复杂对象直接传进子进程的场景下fork更省事。fork启动进程更快spawn本质上相当于fork 通过execv见 POSIX exec 规范额外启动一个不带锁的全新 Python 进程所以 Python 文档明确提示它更慢、开销更大。但请注意我们使用多进程是为了加速那些需要数分钟甚至数小时的计算此时spawn的额外开销在整体耗时面前几乎可以忽略更重要的是它确实能与多线程库一起正常工作。spawn需要代码可被 importspawn启动的是一个全新进程因此被调度的代码必须能被 import反观fork没有这个要求。具体到实践上使用spawn时相关代码不应写在全局作用域里例如 Jupyter notebook 或纯脚本的顶层而是像上文示例那样把真正要并发的逻辑包在函数中并从__main__分支调用。对规范的项目工程这通常不是问题但在 notebook 中快速实验时可能因此失败——这也正是官方示例刻意使用if __name__ __main__:包裹setup()/main()的原因。实践建议把上面的分析收敛成几条可直接执行的经验优先依赖 Polars 自身线程先确认瓶颈不在 Polars 上可用select并行表达式、group_by、惰性引擎的查询优化等手段提速见 优化章节避免盲目叠加多进程。只在「其它单线程库是瓶颈」时引入多进程且用get_context(spawn)显式指定启动方法不要依赖平台默认。把多进程代码收进函数并从__main__启动保证spawn下代码可 import不要在主流程顶层直接创建进程。不要用fork配合任何多线程库forkserver是fork的折中替代其服务器进程本身不加载目标代码、可预加载库但最稳妥、跨平台的选择始终是spawn。若需要限制 Polars 自身的并行度以避免线程池被多进程放大可通过环境变量POLARS_MAX_THREADS或 py-polars/src/polars/meta/thread_pool.py 暴露的 API 控制线程池大小。参考链接Python multiprocessing 文档Polars issue #3144本文示例出处POSIX fork 规范Its faster to fork than to spawn (pythonspeed.com)Python forkserver preloading 讨论【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价