Polars库sink_parquet报错:标准引擎暂不支持该方法
问题描述
我用Polars从大数据集提取特征后,调用sink_parquet写入Parquet文件时触发报错,提示sink_parquet not yet supported in standard engine. Use 'collect().write_parquet()'。但数据集太大,collect()会直接耗尽内存,想知道原因和可行的解决办法。
代码示例
import ipaddress import numpy as np import polars as pl def extract_session_features(sessions: pl.LazyFrame) -> pl.LazyFrame: return ( sessions.with_columns( (pl.col("dpkts") + pl.col("spkts")).alias("total_packets"), (pl.col("dbytes") + pl.col("sbytes")).alias("total_bytes"), (pl.col("dpkts") / pl.col("spkts")).alias("bytes_ratio"), (pl.col("dbytes") / pl.col("sbytes")).alias("packets_ratio"), (pl.col("spkts") / pl.col("dur")).alias("sent_packets_rate"), (pl.col("dpkts") / pl.col("dur")).alias("received_packets_rate"), (pl.col("sbytes") / pl.col("dur")).alias("sent_bytes_rate"), (pl.col("dbytes") / pl.col("dur")).alias("received_bytes_rate"), (pl.col("sbytes") / pl.col("spkts")).alias("mean_pkt_sent_size"), (pl.col("dbytes") / pl.col("dpkts")).alias("mean_pkt_recv_size"), ( pl.col("Timestamp") .diff() .dt.total_seconds() .fill_null(0) .over("ID") .alias("time_since_last_session"), ), ) .with_columns( pl.when(pl.col("^.*(_ratio|_rate).*$").is_infinite()) .then(-1) .otherwise(pl.col("^.*(_ratio|_rate).*$")) .name.keep() ) .fill_nan(-1) ) filtered_sessions = pl.scan_parquet("./processed_merge_file_filtered.parquet") print(filtered_sessions.head().collect(streaming=True)) sessions_features = extract_session_features(filtered_sessions) sessions_features.sink_parquet("./sessions_features")
报错信息
thread '<unnamed>' panicked at /home/runner/work/polars/polars/crates/polars-lazy/src/physical_plan/planner/lp.rs:153:28: sink_parquet not yet supported in standard engine. Use 'collect().write_parquet()' note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace Traceback (most recent call last): File "/home/cpinon/Documentos/Project/SourceCode/project/data/01_raw/sessions/test_feature_engineering.py", line 50, in <module> sessions_features.sink_parquet("./sessions_features") File "/home/cpinon/Documentos/Project/SourceCode/project/.venv/lib/python3.9/site-packages/polars/lazyframe/frame.py", line 1895, in sink_parquet return lf.sink_parquet( pyo3_runtime.PanicException: sink_parquet not yet supported in standard engine. Use 'collect().write_parquet()'
原因分析
报错的核心原因是你的查询中包含了窗口函数(over("ID")):Polars的标准执行引擎无法在流式写入模式下处理需要跨数据行关联计算的窗口操作,而sink_parquet本质是流式输出,因此标准引擎直接抛出了不支持的提示。
可行解决方案
方案1:使用Polars流媒体引擎(推荐)
Polars的流媒体引擎专门针对大内存数据集设计,支持在流式模式下处理窗口函数。只需在调用sink_parquet时指定engine="streaming"即可,注意需要Polars 0.19.0及以上版本:
sessions_features.sink_parquet("./sessions_features", engine="streaming")
如果遇到窗口函数的流式处理限制,可以尝试先按ID和时间戳排序(确保同ID的数据连续),再执行窗口计算:
def extract_session_features(sessions: pl.LazyFrame) -> pl.LazyFrame: return ( sessions .sort("ID", "Timestamp") # 先按ID和时间戳排序,确保同ID数据连续 .with_columns( (pl.col("dpkts") + pl.col("spkts")).alias("total_packets"), (pl.col("dbytes") + pl.col("sbytes")).alias("total_bytes"), (pl.col("dpkts") / pl.col("spkts")).alias("bytes_ratio"), (pl.col("dbytes") / pl.col("sbytes")).alias("packets_ratio"), (pl.col("spkts") / pl.col("dur")).alias("sent_packets_rate"), (pl.col("dpkts") / pl.col("dur")).alias("received_packets_rate"), (pl.col("sbytes") / pl.col("dur")).alias("sent_bytes_rate"), (pl.col("dbytes") / pl.col("dur")).alias("received_bytes_rate"), (pl.col("sbytes") / pl.col("spkts")).alias("mean_pkt_sent_size"), (pl.col("dbytes") / pl.col("dpkts")).alias("mean_pkt_recv_size"), ( pl.col("Timestamp") .diff() .dt.total_seconds() .fill_null(0) .over("ID") .alias("time_since_last_session"), ), ) .with_columns( pl.when(pl.col("^.*(_ratio|_rate).*$").is_infinite()) .then(-1) .otherwise(pl.col("^.*(_ratio|_rate).*$")) .name.keep() ) .fill_nan(-1) ) # 调用流式写入 sessions_features.sink_parquet("./sessions_features", engine="streaming")
方案2:分批次处理数据集
如果流媒体引擎仍有问题,可以将数据集拆分为多个小批次,逐个处理后写入独立的Parquet文件,最后合并所有文件:
# 定义批次大小(根据你的内存情况调整,比如100万行) chunk_size = 1_000_000 # 分批次读取、处理并写入 for i, chunk in enumerate(filtered_sessions.iterate_chunks(chunk_size)): processed_chunk = extract_session_features(chunk) # 写入分块文件 processed_chunk.write_parquet(f"./sessions_features_chunk_{i}.parquet") # 合并所有分块文件为最终结果 final_lf = pl.scan_parquet("./sessions_features_chunk_*.parquet") final_lf.sink_parquet("./sessions_features_final.parquet", engine="streaming")
方案3:调整窗口函数实现(针对性优化)
如果time_since_last_session是唯一的窗口操作,可以尝试用分组后排序的方式替代窗口函数,降低流式处理的复杂度:
def extract_session_features(sessions: pl.LazyFrame) -> pl.LazyFrame: return ( sessions .sort("ID", "Timestamp") .group_by("ID") .agg( pl.exclude("Timestamp"), pl.col("Timestamp").diff().dt.total_seconds().fill_null(0).alias("time_since_last_session") ) .explode(pl.exclude("ID")) .with_columns( (pl.col("dpkts") + pl.col("spkts")).alias("total_packets"), (pl.col("dbytes") + pl.col("sbytes")).alias("total_bytes"), (pl.col("dpkts") / pl.col("spkts")).alias("bytes_ratio"), (pl.col("dbytes") / pl.col("sbytes")).alias("packets_ratio"), (pl.col("spkts") / pl.col("dur")).alias("sent_packets_rate"), (pl.col("dpkts") / pl.col("dur")).alias("received_packets_rate"), (pl.col("sbytes") / pl.col("dur")).alias("sent_bytes_rate"), (pl.col("dbytes") / pl.col("dur")).alias("received_bytes_rate"), (pl.col("sbytes") / pl.col("spkts")).alias("mean_pkt_sent_size"), (pl.col("dbytes") / pl.col("dpkts")).alias("mean_pkt_recv_size"), ) .with_columns( pl.when(pl.col("^.*(_ratio|_rate).*$").is_infinite()) .then(-1) .otherwise(pl.col("^.*(_ratio|_rate).*$")) .name.keep() ) .fill_nan(-1) )
内容的提问来源于stack exchange,提问作者Camilo Piñón
相关产品推荐
相关产品推荐

