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

Python Polars处理30GB文件OOM问题的优化方案咨询

Polars处理30GB时间序列生成OHLC的内存优化方案

问题背景

我需要将原始时间序列数据聚合为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 06:42:18