Python Polars处理30GB文件OOM问题的优化方案咨询
问题背景
我需要将原始时间序列数据聚合为1分钟粒度的OHLC数据,使用Python-Polars编写了代码,执行计划如下:
SELECT [col("parsed_timestamp"), col("label1"), col("open"), col("high"), col("low"), col("close")] FROM AGGREGATE [ col("price").first().alias("open"), col("price").max().alias("high"), col("price").min().alias("low"), col("price").last().alias("close") ] BY [col("label1")] FROM WITH_COLUMNS: [col("timestamp").strict_cast(Datetime(Microseconds, None)).alias("parsed_timestamp")] Csv SCAN [ticker.csv] PROJECT */11 COLUMNS
我使用了scan_csv和collect(streaming=True),同时调用了.group_by_dynamic("parsed_timestamp", every="1m", group_by=["label1"])(该操作未在执行计划中显示)。但处理30GB文件时,代码仍占用30GB内存。数据已按时间戳升序排列。
我编写了一个简单程序,逐行遍历数据并更新内存中的K线,完成后导出,几乎不占用额外内存,Polars能否实现类似效果?
更新后的代码如下:
input_file = "input.csv" output_file = "output.csv" query = ( pl.scan_csv( input_file, has_header=False, new_columns=[ "label1", "timestamp", "price", ], ) .with_columns( pl.col("timestamp").cast(pl.Datetime).alias("parsed_timestamp"), ) .group_by_dynamic("parsed_timestamp", every="1m", group_by=["label1"]) .agg( pl.col("price").first().alias("open"), pl.col("price").max().alias("high"), pl.col("price").min().alias("low"), pl.col("price").last().alias("close"), ) .select( [ "parsed_timestamp", "label1", "open", "high", "low", "close", ] ) ) query.collect(streaming=True).write_csv(output_file)
优化方案
1. 仅读取需要的列
当前执行计划显示PROJECT */11 COLUMNS,说明加载了所有11列,但实际只用到label1、timestamp、price三列。在scan_csv中添加columns参数,只读取必要列,直接减少内存占用:
pl.scan_csv( input_file, has_header=False, new_columns=["label1", "timestamp", "price"], columns=[0, 1, 2], # 对应CSV中这三列的索引,根据实际位置调整 )
2. 提前指定列类型,避免类型推断
Polars默认会自动推断列类型,这个过程会占用额外内存。直接指定列类型可以跳过推断,同时使用更紧凑的类型(比如用Float32替代默认的Float64,如果精度允许):
pl.scan_csv( input_file, has_header=False, new_columns=["label1", "timestamp", "price"], columns=[0, 1, 2], dtypes={ "label1": pl.Categorical, # 如果label1是有限枚举值,用Categorical更节省内存 "timestamp": pl.Datetime, "price": pl.Float32, }, )
3. 用sink_csv替代collect + write_csv
当前代码先通过collect(streaming=True)把结果加载到内存,再写入文件。改用sink_csv可以直接流式输出结果,无需把整个数据集存入内存:
query.sink_csv(output_file, streaming=True)
4. 利用数据排序特性优化group_by_dynamic
因为数据已按时间戳升序排列,可以给group_by_dynamic添加maintain_order=True参数,让Polars利用排序后的结构更高效地分块处理,减少内存开销:
.group_by_dynamic( "parsed_timestamp", every="1m", group_by=["label1"], maintain_order=True )
5. 调整流式处理的块大小
可以通过Polars配置调整流式处理的块大小,避免单块数据占用过多内存:
pl.Config.set_streaming_chunk_size(1024 * 1024 * 64) # 设置为64MB,根据内存情况调整
Polars能否实现低内存逐行式处理?
可以。通过上述优化(尤其是仅读取必要列、指定紧凑类型、流式写入),Polars的处理模式会接近你手动逐行处理的内存占用水平。group_by_dynamic本身就是为时间序列聚合优化的操作,结合流式处理后,Polars会分块读取数据、逐块聚合,聚合完成的块直接写入输出,不会在内存中保留完整的原始数据或结果集,内存占用可以控制在很低的水平。
内容的提问来源于stack exchange,提问作者Ivan

