Polars Lazy DataFrame执行collect(streaming=True)时卡顿无响应求助
Polars流式处理1.56亿行捐款数据时
collect(streaming=True)卡顿问题 我正在用Python+Polars处理2010-2022年公民向政治候选人的个人捐款数据,通过scan_csv加载多份文本文件生成了含1.56亿行、17列的Lazy DataFrame。为后续做字符串匹配关联数据,执行了字符串转小写、正则提取姓名、数据清洗等操作后,调用df.collect(streaming=True)时出现异常:初始CPU、内存及磁盘写入占用下降,数分钟后资源占用回升,代码陷入卡顿或循环,无进展。完整代码如下:
import polars as pl # Define the schema for the data schema = { "Cycle": pl.Int16, "FECTransID": pl.Int64, "ContribID": pl.Utf8, "Contrib": pl.Utf8, "RecipID": pl.Utf8, "Orgname": pl.Utf8, "UltOrg": pl.Utf8, "RealCode": pl.Utf8, "Date": pl.Date, "Amount": pl.Int32, "Street": pl.Utf8, "City": pl.Utf8, "State": pl.Categorical, "Zip": pl.Int32, "RecipCode": pl.Utf8, "Type": pl.Utf8, "CmteID": pl.Utf8, "OtherID": pl.Utf8, "Gender": pl.Categorical, "Microfilm": pl.Utf8, "Occupation": pl.Utf8, "Employer": pl.Utf8, "Source": pl.Utf8 } # Load data with specified schema df = pl.scan_csv( r'C:\Path\To\Your\File\indivs*.txt', separator=',', quote_char='|', schema=schema, encoding='utf8-lossy', ignore_errors=True ).drop(["Microfilm", "OtherID", "Source", "Date", "Street"]) # List of string columns to convert to lowercase string_columns = [ "ContribID", "Contrib", "RecipID", "Orgname", "UltOrg", "RealCode", "City", "RecipCode", "Type", "CmteID", "Occupation", "Employer" ] # Convert string columns to lowercase df = df.with_columns([ pl.col(col).str.to_lowercase().alias(col) for col in string_columns ]) # Collect column names and create a renaming map to lowercase names rename_map = {col: col.lower() for col in df.collect_schema().names()} # Rename the columns df = df.rename(rename_map) # Add a unique identifier column df = df.with_columns(unique_id=pl.int_range(pl.len())) # Define regex patterns for extracting name components regex_last = r'([^,]+)' regex_middle = r'\b([A-Za-z])\b\s*$' regex_first = r',\s+([A-Za-z\s]+)' # Extract name components and clean up first name df = df.with_columns([ pl.col("contrib").str.extract(regex_last, group_index=1).alias("last"), pl.col("contrib").str.extract(regex_middle, group_index=1).alias("middle"), pl.col("contrib").str.extract(regex_first, group_index=1).alias("first") ]).with_columns( pl.col("first").str.replace(r"(\s+[A-Za-z]\.?)\s*$", "").alias("first") ) # Concatenate names into a full name df = df.with_columns( pl.when(pl.col("middle").is_not_null()) .then(pl.col("first") + " " + pl.col("middle") + ". " + pl.col("last")) .otherwise(pl.col("first") + " " + pl.col("last")) .alias("full") ) def remove_all_non_alphanumeric(col): # Remove non-alphanumeric characters from the column return col.str.replace_all("[^a-zA-Z0-9]", "") # Apply the cleaning function to specific columns df = df.with_columns([ remove_all_non_alphanumeric(pl.col("employer")).alias("employer"), remove_all_non_alphanumeric(pl.col("zip")).alias("zip"), remove_all_non_alphanumeric(pl.col("city")).alias("city"), remove_all_non_alphanumeric(pl.col("ultorg")).alias("ultorg"), remove_all_non_alphanumeric(pl.col("orgname")).alias("orgname"), ]) # Collect the final DataFrame, enabling streaming to handle large data df = df.collect(streaming=True)
问题排查与优化方案
1. 消除不必要的Schema预计算
调用df.collect_schema().names()会触发额外的Schema解析开销,直接基于已定义的Schema或当前LazyFrame的列名生成小写映射即可:
# 替换原rename_map生成逻辑 df = df.drop(["Microfilm", "OtherID", "Source", "Date", "Street"]) rename_map = {col: col.lower() for col in df.columns} df = df.rename(rename_map)
2. 优化正则表达式性能
正则操作是字符串处理的性能瓶颈,简化正则逻辑并合并操作:
- 给
regex_last添加开头锚定符^,避免全局扫描 - 合并姓名提取与清洗的多次
with_columns调用,减少查询计划节点
# 合并姓名提取与清洗步骤 df = df.with_columns([ pl.col("contrib").str.extract(r'^([^,]+)', 1).alias("last"), pl.col("contrib").str.extract(r'\b([A-Za-z])\b\s*$', 1).alias("middle"), pl.col("contrib").str.extract(r',\s*([A-Za-z\s]+?)\s*(?:[A-Za-z]\.?)*$', 1).alias("first") ])
3. 移除数值列的无效字符串操作
zip列定义为pl.Int32类型,本身不存在非字母数字字符,无需执行清洗:
# 去掉zip列的清洗操作 df = df.with_columns([ remove_all_non_alphanumeric(pl.col("employer")).alias("employer"), remove_all_non_alphanumeric(pl.col("city")).alias("city"), remove_all_non_alphanumeric(pl.col("ultorg")).alias("ultorg"), remove_all_non_alphanumeric(pl.col("orgname")).alias("orgname"), ])
4. 调整流式处理批次参数
手动设置streaming_buffer_size控制批次大小,避免内存过载:
# 根据本机内存调整,示例设置为1GB批次 df = df.collect(streaming=True, streaming_buffer_size=1024*1024*1024)
5. 拆分查询为多阶段持久化
将复杂处理拆分为多个阶段,用sink_parquet持久化中间结果,避免重复计算:
# 阶段1:加载、转小写、重命名,持久化 df_stage1 = pl.scan_csv( r'C:\Path\To\Your\File\indivs*.txt', separator=',', quote_char='|', schema=schema, encoding='utf8-lossy', ignore_errors=True ).drop(["Microfilm", "OtherID", "Source", "Date", "Street"]) \ .with_columns([pl.col(col).str.to_lowercase().alias(col) for col in string_columns]) \ .rename(rename_map) df_stage1.sink_parquet("stage1_data.parquet", streaming=True) # 阶段2:姓名提取与清洗,持久化 df_stage2 = pl.scan_parquet("stage1_data.parquet") \ .with_columns(unique_id=pl.int_range(pl.len())) \ .with_columns([ pl.col("contrib").str.extract(r'^([^,]+)', 1).alias("last"), pl.col("contrib").str.extract(r'\b([A-Za-z])\b\s*$', 1).alias("middle"), pl.col("contrib").str.extract(r',\s*([A-Za-z\s]+?)\s*(?:[A-Za-z]\.?)*$', 1).alias("first") ]) \ .with_columns( pl.when(pl.col("middle").is_not_null()) .then(pl.col("first") + " " + pl.col("middle") + ". " + pl.col("last")) .otherwise(pl.col("first") + " " + pl.col("last")) .alias("full") ) \ .with_columns([ remove_all_non_alphanumeric(pl.col("employer")).alias("employer"), remove_all_non_alphanumeric(pl.col("city")).alias("city"), remove_all_non_alphanumeric(pl.col("ultorg")).alias("ultorg"), remove_all_non_alphanumeric(pl.col("orgname")).alias("orgname"), ]) df_stage2.sink_parquet("stage2_data.parquet", streaming=True) # 阶段3:最终加载结果 df_final = pl.scan_parquet("stage2_data.parquet").collect(streaming=True)
6. 给正则添加超时限制
Polars 0.19+支持正则超时,避免异常格式字符串触发无限回溯:
pl.col("contrib").str.extract(r',\s+([A-Za-z\s]+)', 1, timeout=1000).alias("first")
内容的提问来源于stack exchange,提问作者Kamron Eck
相关产品推荐
相关产品推荐

