十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

Polars IO Plugins 开发指南:用 Python 注册自定义数据源并接入惰性查询优化

Polars IO Plugins 开发指南:用 Python 注册自定义数据源并接入惰性查询优化 Polars IO Plugins 开发指南用 Python 注册自定义数据源并接入惰性查询优化【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polarsPolars 不仅支持通过表达式插件Expression Plugins扩展函数还提供了一套 IO 插件机制允许你把自定义文件格式或数据源注册为查询引擎的“数据源source”从而自动获得投影下推、谓词下推、提前停止early stopping以及流式引擎的支持。本文将基于官方用户指南与仓库源码从零编写一个极简自定义 CSV 扫描器逐步讲解register_io_source的完整用法、四个预定义参数语义、schema 校验与谓词序列化机制并结合 Rust 侧实现说明底层执行原理让你能够为自己项目中的私有数据格式接入 Polars 的 Lazy 查询管线。为什么需要 IO 插件源码解读在 Polars 中除了表达式插件IO 插件允许你注册不同的文件格式作为查询引擎的数据源。设计文档明确指出两点动机数据源可以通过Arrow FFI 实现零拷贝zero-copy数据移动数据源可以在返回前一次性产出大批量chunk数据因此数据交互的“会合点”仅在短暂时刻需要持有 GIL。正因如此目前 IO 插件接口选择经由 Python 暴露即便数据读取由 Rust 完成也只在会合时短暂需要 GIL不会造成明显的锁竞争。你可以在 py-polars/src/polars/io/plugins.py 中看到该公共 API 的实际实现。典型使用场景当你有 Polars 原生不支持的源文件同时又希望获得以下优化能力时就应该使用 IO 插件投影下推projection pushdown只读取查询真正需要的列谓词下推predicate pushdown在数据源侧尽早过滤行提前停止early stopping例如head(2)时读满指定行数即可终止**流式引擎streaming engine**支持。快速开始编写一个极简 CSV 数据源下面我们按官方用户指南docs/source/user-guide/plugins/io_plugins.md中的思路写一个非常简单甚至低效的自定义 CSV source仅供学习使用。首先导入所需模块# Use python for csv parsing. import csv import polars as pl # Used to register a new generator on every instantiation. from polars.io.plugins import register_io_source from typing import Iterator import io第一步解析 SchemaPolars 中每个scan函数都必须能够提供它读取数据的 schema。对这个简单 CSV 解析器我们一律将数据读为pl.String唯一不同的只是字段名与字段数量def parse_schema(csv_str: str) - pl.Schema: first_line csv_str.split(\n)[0] return pl.Schema({k: pl.String for k in first_line.split(,)})用一个小 CSV 串a,b,c\n1,2,3测试 print(parse_schema(a,b,c\n1,2,3)) Schema([(a, String), (b, String), (c, String)])第二步编写数据源source接下来是真正的 source。我们用一个外层函数加一个内层函数外层函数my_scan_csv是用户面对的函数接受文件名或本例中的整个 CSV 字符串以及其他读源所需参数。对 CSV 而言这些参数可能是delimiter、quote_char等。外层函数调用register_io_source它接受一个callable与一个schema。schema 是整个源文件的 Polars schema与投影下推无关即全量列。内层函数是 IO source 的真正实现也可以继续调用 Rust/C 等其他语言实现的读取逻辑。callable 必须接受以下四个预定义参数参数含义约定with_columns被投影的列若被应用reader 必须对这些列做投影predicatePolars 表达式reader 必须据此过滤行支持谓词下推的源可解析表达式跳过行/组n_rows只从源物化 n 行reader 读到n_rows行即可停止batch_size理想的批次大小提示reader 的生成器应尽力按此大小产出对应代码实现def my_scan_csv(csv_str: str) - pl.LazyFrame: schema parse_schema(csv_str) def source_generator( with_columns: list[str] | None, predicate: pl.Expr | None, n_rows: int | None, batch_size: int | None, ) - Iterator[pl.DataFrame]: Generator function that creates the source. This function will be registered as IO source. if batch_size is None: batch_size 100 # Initialize the reader. reader csv.reader(io.StringIO(csv_str), delimiter,) # Skip the header. _ next(reader) # Ensure we dont read more rows than requested from the engine while n_rows is None or n_rows 0: if n_rows is not None: batch_size min(batch_size, n_rows) rows [] for _ in range(batch_size): try: row next(reader) except StopIteration: n_rows 0 break rows.append(row) df pl.from_records(rows, schemaschema, orientrow) n_rows - df.height # If we would make a performant reader, we would not read these # columns at all. if with_columns is not None: df df.select(with_columns) # If the source supports predicate pushdown, the expression can be parsed # to skip rows/groups. if predicate is not None: df df.filter(predicate) yield df return register_io_source(io_sourcesource_generator, schemaschema)注意几个关键实现细节batch_size缺省为 100作为“理想批次大小”的提示而非强制值n_rows会被逐步扣减且每个批次的规模用min(batch_size, n_rows)收敛保证不多读引擎不需要的行收到StopIteration时把n_rows置 0 退出循环数据读完即自然结束迭代。第三步运行速度很慢但能跑测试脚本csv_str1 a,b,c,d 1,2,3,4 9,10,11,2 1,2,3,4 1,122,3,4 print(my_scan_csv(csv_str1).collect()) csv_str2 a,b 1,2 9,10 1,2 1,122 print(my_scan_csv(csv_str2).head(2).collect())控制台输出shape: (4, 4) ┌─────┬─────┬─────┬─────┐ │ a ┆ b ┆ c ┆ d │ │ --- ┆ --- ┆ --- ┆ --- │ │ str ┆ str ┆ str ┆ str │ ╞═════╪═════╪═════╪═════╡ │ 1 ┆ 2 ┆ 3 ┆ 4 │ │ 9 ┆ 10 ┆ 11 ┆ 2 │ │ 1 ┆ 2 ┆ 3 ┆ 4 │ │ 1 ┆ 122 ┆ 3 ┆ 4 │ └─────┴─────┴─────┴─────┘ shape: (2, 2) ┌─────┬─────┬─────┐ │ a ┆ b │ │ --- ┆ --- │ │ str ┆ str │ ╞═════╪═════╡ │ 1 ┆ 2 │ │ 9 ┆ 10 │ └─────┴─────┴─────┘第二个查询中head(2)之所以只产出两行正是因为 IO 源支持了n_rows提前停止引擎只物化需要的两行数据。深入register_io_source完整签名仓库中 py-polars/src/polars/io/plugins.py 的register_io_source还包含一些官方文档未展开、但实战中非常有用的参数unstable() def register_io_source( io_source: Callable[ [list[str] | None, Expr | None, int | None, int | None], Iterator[DataFrame] ], *, schema: Callable[[], SchemaDict] | SchemaDict, validate_schema: bool False, is_pure: bool False, explain_name: str | None None, explain_detail: str | None None, ) - LazyFrame:io_source即上文四参数签名with_columns,predicate,n_rows,batch_size返回产出 DataFrame 的迭代器/生成器的函数。schemaSchema 或返回 Schema 的函数表示投影下推之前reader 将产出的完整 schema。validate_schema是否让引擎校验产出的批次与给定 schema 一致。不一致属于实现错误会引发难以排查的 bug因此建议开启。仓库测试 py-polars/tests/unit/io/test_io_plugin.py 中的test_defer_validate_true/test_defer_validate_false证实开启时类型不符会抛出pl.exceptions.SchemaError关闭时引擎按给定 schema 解释数据浮点值不会报错。is_pureIO source 是否“纯”。纯数据源在 LazyFrame 计划中重复出现时优化器可以对它做去重de-duplicate。explain_name/explain_detailexplain()查询计划中该 scan 节点的自定义标签与单行细节说明分别渲染为PYTHON[name] SCAN ...与其下的INFO:行。此外该文件还定义了_defer延迟执行工具将“无参函数产出 DataFrame”推迟到 collect 时才运行它内部同样基于register_io_source实现并对with_columns/predicate/n_rows用 Lazy 算子select/filter/limit在收集前应用。底层调用链谓词序列化与 Python Scan 节点register_io_source的核心落在 py-polars/src/polars/lazyframe/frame.py 的LazyFrame._scan_python_function。它依据 schema 形态选择三条路径构造底层逻辑计划普通Mappingschema →PyLazyFrame.scan_from_python_function_pl_schema(...)Arrowpa.Schema→scan_from_python_function_arrow_schema(...)可调用 schema 函数 →scan_from_python_function_schema_function(...)。这些接口最终会生成计划中的 Python scan 节点Rust 侧对应 crates/polars-plan/src/plans/python/source.rs 等实际执行时由 crates/polars-mem-engine/src/executors/scan/python_scan.rs 驱动。因此 IO 插件天然进入 Polars 的优化与执行管线这就是它能获得下推/流式能力的原因。一个容易忽略的机制在 py-polars/src/polars/io/plugins.py 的wrap函数中引擎侧传入的predicate实际是序列化后的字节Python 侧通过pl.Expr.deserialize(predicate)还原为Expr再交给你的生成器。若反序列化失败例如谓词包含 IO 插件无法解析的表达式插件会把parsed_predicate_success置为 false 并告知引擎“过滤由 Polars 侧兜底处理”。设置环境变量POLARS_VERBOSE时会在 stderr 打印形如failed parsing IO plugin expression的提示方便排查。仓库测试 py-polars/tests/unit/io/test_io_plugin.py 中就有对应的回归用例test_io_plugin_predicate_no_serialization_21130str.json_path_match(...)这类表达式无法序列化给插件时引擎自行完成过滤仍返回正确结果test_datetime_io_predicate_pushdown_21790验证pl.col(timestamp) cutoff的谓词确实被推送给了 source记录回调收到的predicate并断言其内容test_io_plugin_reentrant_deadlock单线程POLARS_MAX_THREADS1下 IO source 内部再次调用register_io_source(...).collect()也不会死锁。让 IO 插件配合explain可读的查询计划IO 插件的查询节点同样会出现在explain()输出中。测试test_io_plugin_custom_explain见 py-polars/tests/unit/io/test_io_plugin.py验证了默认标签为PYTHON SCAN当传入explain_nameleft、explain_detailleft detail后计划中会出现PYTHON[left] SCAN ... PROJECT */1 COLUMNS INFO: left detail这对于区分同一查询中多个不同 IO source例如 JOIN 或 UNION 两侧或记录来源配置摘要非常实用。完整生命周期里的实践要点结合 py-polars/tests/unit/io/test_io_plugin.py 中其他用例可以提炼出实现健壮 IO 插件的几条通用准则不要漏掉边界批次生成器应处理“最后一批不足batch_size”与“源读完”两种情况。test_scan_lines中的字节流逐行扫描示例展示了如何在remaining_rows 0时正确 break避免产出空批次或死循环。空源也要有正确 schematest_empty_iterator_io_plugin表明即使生成器不产出任何 DataFrame只要 schema 给定collect()后仍能得到零行且 schema 正确的 DataFrame。投影与重排列顺序test_reordered_columns_22731验证了select(b, a)、with_row_index()、with_columns(b, a)等场景在validate_schemaTrue/False下都能得到与原生 DataFrame 一致的结果——你自己的实现应始终先应用with_columns再应用predicate并保证产出列与 schema 对应。支持的 dtype引擎会按给定 schema 组装批次test_io_plugin_categorical_24172、test_io_plugin_object_dtype_25740分别验证了Categorical与Object类型数据经过 IO 插件后也能保持一致语义多 chunk 数据需注意 categorical 的字典合并。用 Rust 实现 IO source 的参考官方文档推荐参阅 pyo3-polars 仓库中的 io_plugin 示例——该示例在当前仓库内也有完整副本pyo3-polars/example/io_plugin/io_plugin/src/lib.rs。示例中的RandomSource是一个用#[pyclass]暴露给 Python 的随机数据源类实现了典型的“pushdown 友好”接口schema()返回PySchema字段来自各采样器列try_set_predicate(predicate: PyExpr)接收从 Python 侧下推来的谓词set_with_columns(columns: VecString)把列名解析为 schema 索引仅在需要时采样从而“避免不必要的采样”即投影下推落到 reader 内部next()每次按size_hint与剩余n_rows的最小值产出批次切片下推并在批次内应用谓词过滤注释指出对某些源可以在更底层应用。其 Python 侧用法与本文 CSV 示例一致把该 Rust 类包进一个返回Iterator[pl.DataFrame]的生成器再调用register_io_source注册即可。换言之IO 插件的“前后端”分工是Python 负责协议与调度重活解析、过滤、采样可以全部下沉到 Rust这正是 IO 插件适合对接私有高性能读路径的原因。更多资料插件机制总览docs/source/user-guide/plugins/index.md表达式插件指南docs/source/user-guide/plugins/expr_plugins.mdregister_io_source实现py-polars/src/polars/io/plugins.pyIO 插件测试用例py-polars/tests/unit/io/test_io_plugin.pyRust 实现参考发行版示例pyo3-polars/example/io_plugin/io_plugin/src/lib.rs【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表