如何用Polars LazyFrame替代Pivot实现时间戳用户值的高效求和?
问题描述
我有一组带时间戳(timestamp)的事件数据,需要为每个时间戳计算每个用户名(username)对应kudos字段的最后值之和。使用DataFrame的透视表(pivot)可以实现该需求,但当唯一用户名达数千级、事件数据达千万级时,透视表会耗尽内存。希望改用LazyFrame处理,但LazyFrame不支持pivot操作,求高效实现方式(需利用延迟求值,避免生成巨大的稀疏透视表)。
示例输入
import polars as pl df = pl.from_repr(""" ┌────────────┬──────────┬───────┐ │ timestamp ┆ username ┆ kudos │ │ --- ┆ --- ┆ --- │ │ i64 ┆ str ┆ i64 │ ╞════════════╪══════════╪═══════╡ │ 1690886106 ┆ ABC ┆ 123 │ │ 1690886107 ┆ DEF ┆ 10 │ │ 1690886110 ┆ DEF ┆ 12 │ │ 1690886210 ┆ GIH ┆ 0 │ └────────────┴──────────┴───────┘ """)
原DataFrame pivot实现(仅适用于小数据量)
( df.pivot( on="username", index="timestamp", values=["kudos"], aggregate_function="last", ) .select(pl.all().forward_fill()) .fill_null(strategy="zero") .select(pl.col("timestamp"), pl.sum_horizontal(df["username"].unique().to_list()).alias("sum")) )
预期结果
shape: (4, 2) ┌────────────┬─────┐ │ timestamp ┆ sum │ │ --- ┆ --- │ │ i64 ┆ i64 │ ╞════════════╪═════╡ │ 1690886106 ┆ 123 │ │ 1690886107 ┆ 133 │ │ 1690886110 ┆ 135 │ │ 1690886210 ┆ 135 │ └────────────┴─────┘
LazyFrame高效实现方案
核心思路是避免生成宽表,转而通过分组计算每个用户的最新kudos状态,再结合全局时间轴做聚合求和,全程利用LazyFrame的延迟求值特性,内存占用远低于透视表方案。
方案1:基于笛卡尔积的简洁实现
适合用户数量不是极端庞大的场景,逻辑直观:
import polars as pl # 构造LazyFrame(实际场景可从文件扫描:pl.scan_csv("large_data.csv")) lf = pl.scan_repr(""" ┌────────────┬──────────┬───────┐ │ timestamp ┆ username ┆ kudos │ │ --- ┆ --- ┆ --- │ │ i64 ┆ str ┆ i64 │ ╞════════════╪══════════╪═══════╡ │ 1690886106 ┆ ABC ┆ 123 │ │ 1690886107 ┆ DEF ┆ 10 │ │ 1690886110 ┆ DEF ┆ 12 │ │ 1690886210 ┆ GIH ┆ 0 │ └────────────┴──────────┴───────┘ """) result = ( # 1. 获取所有唯一时间戳并排序,作为全局时间轴 lf.select("timestamp").unique().sort("timestamp") # 2. 和每个用户的历史kudos数据做笛卡尔积,覆盖所有时间点-用户组合 .join(lf.group_by("username").agg(pl.col("timestamp"), pl.col("kudos")), how="cross") # 3. 筛选每个用户在当前时间点之前的最后一条kudos记录 .filter(pl.col("timestamp_right") <= pl.col("timestamp")) .sort("timestamp", "timestamp_right") .group_by("timestamp", "username") .last() # 4. 按时间戳聚合求和,空值自动视为0(用户未出现时无贡献) .group_by("timestamp") .agg(pl.col("kudos").sum().alias("sum")) .sort("timestamp") # 触发计算(延迟求值到最后一步) .collect() ) print(result)
方案2:无笛卡尔积的内存友好实现
适合用户数量极多的场景,避免笛卡尔积带来的临时数据膨胀:
result = ( # 1. 按用户分组,保留每个用户的时间戳和kudos记录 lf.group_by("username") .agg(pl.col("timestamp"), pl.col("kudos")) # 2. 展开用户的历史记录,和全局时间轴做外连接,补全所有时间点 .explode(["timestamp", "kudos"]) .join(lf.select("timestamp").unique().sort("timestamp"), on="timestamp", how="outer") # 3. 按用户排序后,向前填充kudos值(无新数据时保持最后状态),未出现的用户填充0 .sort("username", "timestamp") .group_by("username") .agg( pl.col("timestamp"), pl.col("kudos").forward_fill().fill_null(0) ) .explode(["timestamp", "kudos"]) # 4. 按时间戳聚合求和 .group_by("timestamp") .agg(pl.col("kudos").sum().alias("sum")) .sort("timestamp") .collect() ) print(result)
输出验证
两种方案的输出均与预期结果一致,且全程使用LazyFrame的延迟求值,仅在最后collect()时才计算,内存占用可控,适合千万级事件数据、数千级用户的场景。
内容的提问来源于stack exchange,提问作者sougonde
相关产品推荐
相关产品推荐

