Polars处理大数据集致Python内核崩溃问题求助
问题分析
你遇到的内核崩溃、内存泄漏问题,核心原因是**np.vectorize与Polars map_batches结合时的内存管理冲突**:
np.vectorize本质是Python循环的包装,在Polars批量处理场景下,会频繁在Polars的Arrow内存和NumPy数组之间做数据拷贝;- 叉乘后的DataFrame数据量极大(比如两个1万行的DF叉乘会生成1亿行),这种高频拷贝会导致内存暴增,甚至触发泄漏。
而你用循环处理更大的合成数据时,内存是逐行/逐块释放的,不会出现这种批量拷贝的累积问题。
解决方案
1. 用map_elements替代map_batches+np.vectorize
直接针对单个struct元素逐行计算,避免批量转NumPy数组的内存开销,Polars会自动管理内存:
import polars as pl import jellyfish cross_df.with_columns( pl.struct(["title", "title_right"]) .map_elements(lambda x: jellyfish.jaro_winkler_similarity(x["title"], x["title_right"])) .alias("JW_similarity") )
2. 自定义批量处理的Polars UDF
直接处理Arrow列转成的Python列表,批量计算相似度,比逐行处理更高效:
import polars as pl from polars import Expr import jellyfish from typing import List def jaro_winkler_expr(left: Expr, right: Expr) -> Expr: def _batch_calc(left_vals: List[str], right_vals: List[str]) -> List[float]: return [jellyfish.jaro_winkler_similarity(l, r) for l, r in zip(left_vals, right_vals)] return pl.struct([left, right]).map_batches( lambda s: pl.Series(_batch_calc(s.struct.field(left.name).to_list(), s.struct.field(right.name).to_list())) ) # 调用自定义UDF cross_df.with_columns( jaro_winkler_expr(pl.col("title"), pl.col("title_right")).alias("JW_similarity") )
3. 分块处理超大规模DataFrame
如果叉乘后数据量突破内存上限(比如10亿行),手动分块处理后合并:
chunk_size = 1_000_000 result_chunks = [] for i in range(0, len(cross_df), chunk_size): chunk = cross_df.slice(i, chunk_size).with_columns( pl.struct(["title", "title_right"]) .map_elements(lambda x: jellyfish.jaro_winkler_similarity(x["title"], x["title_right"])) .alias("JW_similarity") ) result_chunks.append(chunk) final_result = pl.concat(result_chunks)
额外优化建议
- 先过滤再叉乘:如果不是必须全量叉乘,可先按
author分组,只在同组内做叉乘,大幅减少数据量; - 替换计算库:jellyfish是纯Python实现,性能有限,数据量极大时可改用Rust编写的相似性工具(如
polars-similarity扩展),或用Numba加速jellyfish的计算逻辑。
内容的提问来源于stack exchange,提问作者Markos Kapes
相关产品推荐
相关产品推荐

