You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 16:14:53