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
相关产品推荐
相关产品推荐

