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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 13:11:00