ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Polars 惰性查询复用(Multiplexing):用 pl.collect_all 消除重复子计划的重复计算

Polars 惰性查询复用(Multiplexing):用 pl.collect_all 消除重复子计划的重复计算 Polars 惰性查询复用Multiplexing用 pl.collect_all 消除重复子计划的重复计算【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars在 Polars 的LazyFrame世界里一个常见的认知陷阱是把惰性查询当成已经算好的中间结果存进变量。本指南围绕官方用户手册中的 multiplexing 专题系统讲解查询分叉/多路复用multiplexing这一核心概念为何在 eager API 下随手可用的临时变量 过程式循环在 lazy API 下会导致重复计算与结果不稳定以及如何借助pl.collect_all配合pl.explain_all观察让多个源自同一子计划的查询共享同一次优化与执行。读完本文你将掌握在惰性 API 下安全实现结果复用的三种正确姿势并理解CACHE/SINK_MULTIPLE节点背后的实现原理。背景什么是 MultiplexingMultiplexing多路复用在 Polars 文档中首先出现在 Sources and sinks 一页。在那里它指的是**把一条查询同时写入多个 sink落盘目标**的能力。典型写法如下# 一些昂贵的计算 lf: pl.LazyFrame q1 lf.sink_parquet(..., lazyTrue) # 同时写 parquet q2 lf.sink_ipc(..., lazyTrue) # 同时写 ipc pl.collect_all([q1, q2])在这段代码中q1与q2共享底层的lfsink 的lazyTrue参数使其返回LazyFrame而非立即执行。如果分别collect昂贵的公共上游会被重复计算两次而pl.collect_all([q1, q2])将两个计划合并进一次执行这正是multiplex的含义。本篇 multiplexing.md 则更进一步不限于 sink而是讨论当用过程式编程构造变量暂存、循环切片组合多个LazyFrame时如何理解并利用这一机制。核心矛盾变量里的 LazyFrame 并不是表首先构造一个便于观察的测试数据。手册配套示例脚本位于 docs/source/src/python/user-guide/lazy/multiplexing.py用固定随机种子生成 10 个乱序的互不相同的整数import polars as pl import numpy as np np.random.seed(0) a np.arange(0, 10) np.random.shuffle(a) df pl.DataFrame({n: a}) print(df)之所以刻意打乱顺序注释里写得很明白so that Polars doesnt hit any fast paths for sorted keys——避免 Polars 对已排序键走任何快速路径从而让group_by的输出顺序真正地不确定便于暴露问题。注意group_by之后结果的行序是不保证的。这是 Polars 的既定行为后续所有现象都与这一点相关。Eager直觉正确如果全程使用 eager APIDataFrame直觉写法的结果完全符合预期# A group-by doesnt guarantee order df1 df.group_by(n).len() # Take the lower half and the upper half in a list out [df1.slice(offseti * 5, length5) for i in range(2)] # Assert df1 is equal to the sum of both halves pl.testing.assert_frame_equal(df1, pl.concat(out))关键在第二行注释与第三行注释之间df.group_by(n).len()的结果被立即物化并存入df1。随后对df1做两次slice切前半 5 行与后半 5 行再拼接必然还原出df1本身。虽然group_by输出的行序不稳定但因为所有操作都是即时eager求值顺序已定格在df1里slice与concat都作用于同一份真实数据断言自然通过。Lazy同样写法直接翻车把上面代码机械地翻译成LazyFrame立刻失败lf1 df.lazy().group_by(n).len() out [lf1.slice(offseti * 5, length5).collect() for i in range(2)] pl.testing.assert_frame_equal(lf1.collect(), pl.concat(out))运行时会抛出形如以下的断言错误AssertionError: DataFrames are different (value mismatch for column n) [left]: [9, 2, 0, 5, 3, 1, 7, 8, 6, 4] [right]: [0, 9, 6, 8, 2, 5, 4, 3, 1, 7]左、右两侧明明都源自对同一df做相同group_by结果却顺序不同、无法通过assert_frame_equal。根因变量里存的是查询计划不是结果文档点明了问题的本质lf1并不包含df.lazy().group_by(n).len()的物化结果它持有的只是查询计划query plan。在惰性语义下lf1是一个待执行的蓝图。因此重复计算每次从lf1分叉出一个分支并调用.collect()引擎都会从头重新求值一遍group_by。上述失败示例实际上求值了 2 份此处实为 3 次 collect但共享同一重复模式相互独立的查询计划公共上游被重复执行浪费计算资源。结果不一致若你潜意识里假设输出是稳定的以为lf1是已经算好的表就会踩坑——每次重算的group_by行序并不稳定切片拼接后自然对不上。这正是文档强调的理解 multiplexing 对用过程式编程构造组合LazyFrame至关重要——因为你必须时刻提醒自己变量 计划不是表。两份独立计划的直观呈现用.show_graph()可以可视化这一事实。手册配套代码分别绘制了两个分支的未优化 IR 图对应plan_1/plan_2两段q1 df.lazy().group_by(n).len() q2 q1.slice(offset0, length5) # 第一个分支 show_plan(q2, optimizedFalse)与q1 df.lazy().group_by(n).len() q2 q1.slice(offset5, length5) # 第二个分支 show_plan(q2, optimizedFalse)show_plan只是把LazyFrame.show_graph生成的 PNG 以 base64 内嵌显示的工具函数见 multiplexing.py。观察这两张图可以发现每个分支的图里都各自完整地包含了一份group_by。也就是说逻辑上它们明明共享q1但物理计划层面却是两条互不相干的独立流水线。解法用 pl.collect_all 合并查询计划要绕过上述问题就必须给 Polars 一个机会让它在单次优化与执行过程中同时看到所有分叉的计划。做法是把这些分叉的LazyFrame一起传给pl.collect_alllf1 df.lazy().group_by(n).len() out [lf1.slice(offseti * 5, length5) for i in range(2)] results pl.collect_all([lf1] out) pl.testing.assert_frame_equal(results[0], pl.concat(results[1:]))这里pl.collect_all同时接收了三份计划lf1完整 group-by 结果以及它的两个切片分支。引擎把它们当成一个整体来优化与执行lf1只被计算一次两个切片分支直接消费这份结果返回的results是list[DataFrame]顺序与输入LazyFrame列表一一对应results[0]即lf1的结果results[1:]是两个切片于是concat(results[1:]) results[0]断言通过。用 explain_all 观察共享痕迹如果想眼见为实可以用pl.explain_all查看合并后的物理计划lf1 df.lazy().group_by(n).len() out [lf1.slice(offseti * 5, length5) for i in range(2)] print(pl.explain_all([lf1] out))文档指出pl.explain_all的输出中能够观察到两个关键特征合并后的多个查询被包在一个 SINK_MULTIPLE 求值单元之下优化器识别出多个查询源自同一子计划在公共子计划上插入 CACHE 节点——这与每个分支独立.collect()时图里各画一份 group_by形成鲜明对比。这段可观察现象正是后续 源码印证 部分要展开的内容。API 细节collect_all 与 explain_all 的完整签名collect_all与explain_all均定义在 py-polars/src/polars/functions/lazy.py 中collect_all见 L1970-L2045explain_all见 L2160-L2187。pl.collect_all的关键签名如下pl.collect_all( lazy_frames: Iterable[LazyFrame], # 待收集的 LazyFrame 列表 *, optimizations: QueryOptFlags DEFAULT_QUERY_OPT_FLAGS, engine: EngineType auto, # auto | in-memory | streaming | gpu lazy: bool False, ) - list[DataFrame] | LazyFrame几个值得注意的参数行为lazy_frames一批要一并收集的LazyFrame。返回的DataFrame列表与输入顺序一致。optimizations控制优化 pass 的参数默认取DEFAULT_QUERY_OPT_FLAGS。其中Common Subplan Elimination公共子计划消除正是本页主题的开关——可用pl.QueryOptFlags(comm_subplan_elim...)显式控制。docstring 明确指出Common Subplan Elimination is applied on the combined plan, meaning that diverging queries will run only once.合并计划会应用公共子计划消除使分叉查询只运行一次。engine执行引擎默认auto受Config.set_engine_affinity或环境变量POLARS_ENGINE_AFFINITY影响缺省回落到in-memory也可显式指定in-memory、streaming或gpu。lazyTrue返回一个以后可再收集的LazyFrame仅当所有输入都以 sink 落盘时语义才正确该参数被标记为 unstable。此外 Polars 还提供了异步版本pl.collect_all_asyncL2077-L2156将收集调度进线程池、立刻返回可等待对象支持geventTrue返回gevent.event.AsyncResult适合在 asyncio/gevent 环境中释放事件循环。而pl.explain_all的作用正如其名Explain multiple LazyFrames as if passed tocollect_all——以collect_all的方式合并并展示计划文本让你在不真正执行的情况下检查 CACHE 节点与共享结构。它的optimizations同样可配且整个 API 标有 unstable 标记行为可能在未来版本调整。在 Rust 侧explain_all/collect_all对应的入口位于 crates/polars-python/src/functions/lazy.rs上层引擎分发逻辑见 py-polars/src/polars/lazyframe/engine.py。源码印证CACHE 与 SINK_MULTIPLE 从何而来公共子计划消除CSPE优化器合并计划时识别共享子计划并插入 CACHE这一步由 Polars 优化器中的Common Subplan Elimination模块负责其实现位于 crates/polars-plan/src/plans/optimizer/cse/cspe.rs注册于 crates/polars-plan/src/plans/optimizer/mod.rs。结合源码与测试可以确认几件事实CSPE 针对合并后的多查询生效。这正是collect_all/explain_all的价值只有多个计划被并置优化时优化器才有机会发现跨查询的重复子计划从而算一次、多处引用。显式缓存与自动缓存并存。除了collect_all触发的自动消除外用户也可用LazyFrame.cache()显式要求缓存中间结果。例如回归测试 py-polars/tests/unit/lazyframe/test_cse.py#L1811-L1821 同时用pl.collect_all(frames)与关闭comm_subplan_elim的版本做等价性对比。CSPE 保证结果等价。测试文件中的大量用例如 test_cse.py#L1024-L1163验证开启 CSPE 后pl.concat(pl.collect_all(lfs))的结果与未消除共享时一致。collect_all是否真正共享执行可由explain_all中是否出现 CACHE 节点判断。SINK_MULTIPLE多输出的执行单元在合并计划的文本表示中多个输出被组织在同一个SINK_MULTIPLE ... END SINK_MULTIPLE块下。该格式定义于 crates/polars-plan/src/plans/ir/format.rs 与 tree_format.rs执行期多路复用逻辑由 polars-lazy 的 multiplexed 执行路径承担参见 crates/polars-lazy/src/frame/mod.rs 中与collect_all/ sink 相关的实现。从结构上可以推断SINK_MULTIPLE把共享上游的多个下游输出焊接成一个统一的执行 DAG上游数据只在内存中产生一份经 CACHE 节点被多个输出分支复用。测试用例提供的操作模板仓库内 py-polars/tests/unit/lazyframe/test_collect_all.py 汇集了collect_all的常见用法与边界回归test_collect_allL16-L22两个独立LazyFrame一个是int求和、一个是floats*2求和一次性收集返回结果按输入顺序取出并逐项断言——展示了collect_all也适用于完全不相干的查询并行收集。test_collect_all_groupby_lazy_sink_issue_26296L41-L56把基于同一 group-by 结果派生出的 select 分支与sink_parquet(lazyTrue) 落盘查询一起collect_all——正是本文 sink 复用场景 的最小复现。这些测试从侧面证明collect_all是同时满足结果按序返回共享子计划只算一次支持 sink 与普通分支混用的统一入口。复用场景与落地建议综合文档、配套示例与源码测试pl.collect_all适用的典型复用场景有三类复用场景之一同一结果派生的多个分支即本文主角的过程式构造模式lf expensive_scan_or_groupby(...) # 昂贵的公共子计划 branch1 lf.filter(...) # 多个派生分支 branch2 lf.select(...) branch3 lf.slice(offset0, length100) results pl.collect_all([lf, branch1, branch2, branch3])相比[x.collect() for x in ...]它把一次昂贵计算如扫描、group-by、join从 N 次降为 1 次同时保证了分支间看到的是同一份同一顺序的中间结果规避了非确定性。复用场景之二纯并行收集不相干查询collect_all的 docstring 明确说它可以 run all the computation graphs in parallel or combined。如果多个LazyFrame互不共享子计划传给它等价于一次并行 collectresults pl.collect_all([q_read_csv, q_read_parquet, q_memory])复用场景之三同时落到多个 sink回到 Sources and sinks 中的经典用法——把一份昂贵计算的结果同时以不同格式/分区落盘lf: pl.LazyFrame # 昂贵的查询 q1 lf.sink_parquet(out_a/, lazyTrue) q2 lf.sink_ipc(out_b.ipc, lazyTrue) pl.collect_all([q1, q2])sink_parquet/sink_ipc的lazyTrue让落盘动作本身惰性化并返回LazyFrame再经由collect_all合并底层公共计算只执行一次两个 sink 并行写出。这正是collect_all(..., lazyTrue)参数存在的原因——它返回的仍是一个可继续收集的组合LazyFrame适合与 sink 链式衔接unstable使用时留意警告。何时不必用 collect_all如果各LazyFrame之间没有共享子计划单独collect与collect_all在重复计算层面没有差别仅执行调度/引擎选择可能不同。另外Polars 内部对单个查询哪怕含重复表达式本身也有表达式级别的去重与缓存优化collect_all专门解决的是跨独立查询对象的共享问题。因此最佳实践是把共享上游的派生查询收集到一起提交其余情况按需并行即可。小结本页要点可浓缩为一句话惰性查询变量保存的是计划而非结果重复collect会重复计算且结果顺序不稳定pl.collect_all/pl.explain_all通过跨查询的公共子计划消除CSPE将源自同一子计划的分叉查询合并到单个SINK_MULTIPLE执行单元以CACHE节点共享计算结果实现一次计算、多路复用。回顾全文关键结论关注点eagerDataFramelazy 但各自 collectlazy pl.collect_all变量中保存的物化结果查询计划查询计划提交时合并公共 group-by 执行次数1 次每个分支各 1 次重复1 次CSPE 共享分支间顺序一致性稳定已物化不稳定稳定同一次计算的同一份结果结果直觉正确断言失败断言通过配套示例代码的完整版本可对照 docs/source/src/python/user-guide/lazy/multiplexing.py 阅读想进一步理解扫描与汇出侧的多路复用可继续阅读 Sources and sinks想深挖优化器如何识别共享子计划源码入口在 crates/polars-plan/src/plans/optimizer/cse/cspe.rs。【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表