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

Polars流式处理:基于shift(-1)降采样并写入Parquet的内存优化问题

Polars流式处理:基于shift(-1)降采样并写入Parquet的内存优化问题

首先,我完全理解你的痛点:明明用了scan_parquet和流式引擎,结果内存还是暴涨60GB,这和你预期的O(1)内存处理完全不符。问题出在**shift(-1)这个全局操作上**——让我一步步给你拆解原因,再给出解决办法。

为什么当前代码内存爆炸?

你的过滤逻辑是想保留每个ticker在每个时间桶的最后一条记录(当ticker不变时,ts_bucket变化的那一行;或者ticker变化的最后一行),但用shift(-1)来实现时,Polars需要全局访问下一行数据:

  • 流式模式下,Polars默认的分区策略无法处理这种跨任意行的依赖关系——如果数据没有按ticker分区/排序,Polars无法确定下一行属于哪个ticker,只能把整个数据集加载到内存中计算shift,这就是内存飙到60GB的核心原因。

优化方案:用分组聚合替代全局Shift

你的需求本质上是对每个(ticker, ts_bucket)组取最后一条记录,完全可以用group_by + agg来实现,这完美适配流式处理,内存占用会降到单个ticker的数据量级别(接近O(1))。

优化后的代码

DOWNSAMPLE_NANOS = int(1e11)

# 流式读取Parquet
d = pl.scan_parquet("example_input.parquet")

# 生成时间桶字段
d = d.with_columns(
    (pl.col('epoch_nanos') // DOWNSAMPLE_NANOS).alias('ts_bucket')
)

# 核心优化:按ticker+时间桶分组,取每组最后一行
d = d.group_by(['ticker', 'ts_bucket']).agg(
    # 取每组的最后一条记录,完全等价于你的过滤逻辑
    pl.col('epoch_nanos').last(),
    pl.col('price').last(),
    pl.col('ticker').first()  # 也可以直接用pl.all().last(),按需选择列
)

# 移除时间桶字段(可选)
d = d.drop('ts_bucket')

# 流式写入结果,内存可控
d.sink_parquet("example_output_optimized.parquet", engine='streaming')

为什么这个方案内存友好?

  • group_by(['ticker', 'ts_bucket'])在流式模式下,Polars会自动按ticker分区处理——每个批次只加载单个ticker的数据,处理完该ticker的所有时间桶后就释放内存,不需要全局加载所有10亿行数据。
  • 聚合操作只保留每组的最后一行,数据量直接从10亿行降到10,000个ticker × 每个ticker的时间桶数,内存占用大幅降低。

如果你确实需要用Shift逻辑(复杂场景)

如果你的过滤逻辑比当前更复杂,必须用shift实现,那你需要给Shift操作加上分区限定,让Polars知道只在同一个ticker内计算下一行:

DOWNSAMPLE_NANOS = int(1e11)

d = pl.scan_parquet("example_input.parquet")
d = d.with_columns(
    (pl.col('epoch_nanos') // DOWNSAMPLE_NANOS).alias('ts_bucket')
)

# 关键:用over('ticker')限定Shift只在同一个ticker的范围内执行
d = d.with_columns(
    pl.col('ticker').shift(-1).fill_null('EOF').over('ticker').alias('next_ticker'),
    pl.col('ts_bucket').shift(-1).over('ticker').alias('next_ts_bucket')
)

# 原过滤逻辑,现在只在单个ticker分区内计算
d = d.filter(
    (pl.col('ticker') != pl.col('next_ticker')) | (pl.col('ts_bucket') != pl.col('next_ts_bucket'))
).drop(['next_ticker', 'next_ts_bucket', 'ts_bucket'])

d.sink_parquet("example_output_shift_optimized.parquet", engine='streaming')

这里的over('ticker')告诉Polars:Shift操作只在当前ticker的分区内进行,流式处理可以逐个ticker处理,不需要全局数据,内存占用同样可控。

版本兼容性验证

你使用的Polars 1.32.3已经完全支持流式模式下的group_by和窗口函数over,上述方案都可以正常运行。

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 06:55:32