Polars中如何按分组拆分大型DataFrame为Vec<DataFrame>以适配实时场景?
高效处理Polars实时分组数据的方案
核心思路
直接维护一个以date为键、对应分组DataFrame为值的哈希表(Python字典),新数据到来时仅更新对应分组,避免全量拼接的性能损耗。
具体实现步骤
- 初始化分组字典:将历史数据按
date分组后直接转为字典,每个键对应该日期的完整数据片段。 - 增量更新分组:新数据先按
date分组,遍历每个分组,若日期已存在则用vstack(同结构垂直拼接,比concat更高效)合并,不存在则直接添加到字典。 - 按需聚合计算:需要分析时,可直接遍历字典计算每个分组的聚合结果,再合并为最终结果;也可按需拼接指定分组的DataFrame进行更复杂的分析。
代码示例
1. 初始化历史分组数据
import polars as pl # 模拟历史金融数据 historical_df = pl.DataFrame({ "datetime": pl.date_range(start="2024-01-01", end="2024-01-03", interval="1h"), "price": pl.randn(72), "volume": pl.randint(100, 1000, 72) }).with_columns(pl.col("datetime").dt.date().alias("date")) # 转成以日期为键的分组字典 grouped_data = dict(historical_df.group_by("date"))
2. 处理实时新增数据
# 模拟实时生成的新数据 new_data = pl.DataFrame({ "datetime": pl.date_range(start="2024-01-03", end="2024-01-04", interval="1h"), "price": pl.randn(24), "volume": pl.randint(100, 1000, 24) }).with_columns(pl.col("datetime").dt.date().alias("date")) # 增量更新分组字典 for date, new_group in new_data.group_by("date"): if date in grouped_data: # 用vstack高效拼接同结构数据,避免全量合并 grouped_data[date] = grouped_data[date].vstack(new_group) else: grouped_data[date] = new_group
3. 执行聚合计算
# 计算每个日期的平均价格、总成交量 aggregated_results = [] for date, df in grouped_data.items(): group_agg = df.select( pl.lit(date).alias("date"), pl.col("price").mean().alias("avg_price"), pl.col("volume").sum().alias("total_volume") ) aggregated_results.append(group_agg) # 合并所有分组的聚合结果 final_agg_df = pl.concat(aggregated_results) print(final_agg_df)
注意事项
- 确保新旧数据的Schema完全一致,否则
vstack会报错,可提前用new_data = new_data.cast(historical_df.schema)强制对齐。 - 若
datetime粒度是小时/分钟,需统一分组键的格式(比如按天分组就转成date类型,按小时就转成datetime.truncate("1h"))。 - 超大规模数据场景下,可将非活跃日期的分组数据写入Parquet文件,仅保留近期数据在内存,需要时再读取,进一步降低内存占用。
内容的提问来源于stack exchange,提问作者Hakase
相关产品推荐
相关产品推荐

