如何用Polars懒式流式将API结果写入文件避免内存溢出?
Polars流式追加Parquet+懒加载读取方案
不用每次重新扫描旧文件再拼接,Polars原生支持Parquet追加写入,搭配scan_parquet的懒加载特性,完美满足你每次5万条数据流式写入、无需加载全量内存的需求。
具体实现步骤
1. 首次写入数据
第一次获取到5万条数据后,直接写入Parquet文件:
import polars as pl # 替换成你的API获取逻辑,返回Polars DataFrame df_first_batch = your_api_fetch_function() # 首次写入用overwrite模式,确保文件初始化 df_first_batch.write_parquet("stream_data.parquet", write_mode="overwrite")
2. 后续批次追加
之后每次拿到新的5万条数据,直接用append模式写入,不需要读取旧文件:
df_new_batch = your_api_fetch_function() # 追加到已有Parquet文件,内存只存当前批次数据 df_new_batch.write_parquet("stream_data.parquet", write_mode="append")
3. 懒加载读取全量数据
需要处理全量数据时,用scan_parquet生成LazyFrame,全程不加载全量数据到内存:
# 懒加载扫描Parquet文件,返回LazyFrame lazy_full_data = pl.scan_parquet("stream_data.parquet") # 示例:做过滤、聚合等操作,都是懒执行,直到collect才计算 filtered_result = lazy_full_data.filter(pl.col("status") == "success").group_by("category").agg(pl.sum("amount")).collect()
为什么这个方案比你原来的思路更好
- 省掉了重复读取旧文件的IO成本,每次只写新数据,效率更高
- 内存占用始终只限于当前批次的5万条数据,完全符合流式处理要求
scan_parquet天然支持读取追加后的Parquet文件,不需要额外处理拼接逻辑
注意点
- 每次追加的DataFrame必须列名、数据类型完全一致,否则会触发写入错误
- 如果想按批次拆分文件(比如
batch_1.parquet、batch_2.parquet),可以用pl.scan_parquet("batch_*.parquet")批量扫描,效果和单文件一致,还方便管理 - 确保你的Polars版本在0.15.0以上,这个版本才正式支持Parquet的追加写入
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

