
本文介绍如何使用 Polars 原生表达式高效计算事务流(如按时间戳记录的金额)在指定时间网格(如每分钟整点)上的运行累计和,避免低效的 map_elements 或 apply,全程基于向量化操作。
本文介绍如何使用 polars 原生表达式高效计算事务流(如按时间戳记录的金额)在指定时间网格(如每分钟整点)上的运行累计和,避免低效的 `map_elements` 或 `apply`,全程基于向量化操作。
在时序分析场景中,常需将离散事件(如交易流水)映射到规则时间网格(如每5分钟、每小时),并计算截至每个网格点的累积值(例如“截至 11:03 的总入账额”)。Pandas 用户可能习惯用 apply 配合时间过滤,但在 Polars 中,此类操作应优先采用声明式、向量化的原生表达式(Expression API),而非 Python 回调(如 map_elements),以保障性能与可扩展性。
核心思路是:将原始事件按“首次覆盖的时间网格点”进行归因,再聚合+累加。具体分三步实现:
✅ 步骤 1:对齐事件到目标时间网格(join_asof + forward)
使用 join_asof(..., strategy="forward"),为 df 中每个事件时间查找 dg 中首个 ≥ 该事件时间的网格时间点。这相当于“把这笔交易归入它生效的第一个统计时刻”。
import polars as pl
import datetime
df = pl.DataFrame({
"time": [
datetime.datetime(2025, 2, 2, 11, 1),
datetime.datetime(2025, 2, 2, 11, 2),
datetime.datetime(2025, 2, 2, 11, 3)
],
"amount": [5.0, -1.0, 10.0]
})
dg = pl.DataFrame({
"time": pl.datetime_range(
datetime.datetime(2025, 2, 2, 11, 0),
datetime.datetime(2025, 2, 2, 11, 5),
"1m",
eager=True
)
})
# 关键:前向对齐 → 每个事件绑定到其“生效的最早网格时间”
aligned = df.join_asof(dg, on="time", strategy="forward", coalesce=False)
# 重命名 time_right 为 time,便于后续 join
aligned = aligned.select("amount", pl.col("time_right").alias("time"))
✅ 步骤 2:左连接回时间网格,补全空缺
将对齐后的 amount 映射回 dg,未匹配到事件的网格点 amount 为 null:
mapped = dg.join(aligned, on="time", how="left") # shape: (6, 2) —— 包含 11:00 ~ 11:05 共6个点
✅ 步骤 3:聚合去重 + 累积求和
同一网格时间点可能对应多个事件(如两笔交易均发生在 11:02:30,都归入 11:03),因此需先按 time 分组求和,再计算累积和。注意:为保持结果完整性,需将 null 替换为 0.0(或使用 fill_null(0)):
result = (
mapped
.group_by("time")
.agg(pl.col("amount").sum().fill_null(0))
.sort("time") # 确保时间有序(group_by 不保证顺序)
.with_columns(
pl.col("amount").cum_sum().alias("cum_amount")
)
)
最终输出:
shape: (6, 3) ┌─────────────────────┬────────┬────────────┐ │ time ┆ amount ┆ cum_amount │ │ --- ┆ --- ┆ --- │ │ datetime[μs] ┆ f64 ┆ f64 │ ╞═════════════════════╪════════╪════════════╡ │ 2025-02-02 11:00:00 ┆ 0.0 ┆ 0.0 │ │ 2025-02-02 11:01:00 ┆ 5.0 ┆ 5.0 │ │ 2025-02-02 11:02:00 ┆ -1.0 ┆ 4.0 │ │ 2025-02-02 11:03:00 ┆ 10.0 ┆ 14.0 │ │ 2025-02-02 11:04:00 ┆ 0.0 ┆ 14.0 │ │ 2025-02-02 11:05:00 ┆ 0.0 ┆ 14.0 │ └─────────────────────┴────────┴────────────┘
⚠️ 注意事项与最佳实践
-
时间列必须有序且类型一致:
join_asof要求on列已升序排序(可提前df.sort("time")),且df与dg的time列 dtype 完全相同(推荐datetime[μs])。 -
避免
map_elements:它会触发 Python 解释器循环,丧失 Polars 的零拷贝与并行优势;本方案全程在 Rust 引擎内执行。 -
扩展性提示:若
dg是高频网格(如秒级),而df事件稀疏,此方法仍高效;若需动态窗口(如“过去30分钟累计”),则应改用group_by_dynamic或rolling。 -
空值处理:
.fill_null(0)确保cum_sum()从 0 开始累加,符合业务语义(无交易即无变动)。
通过 join_asof + group_by + cum_sum 三步组合,你能在 Polars 中以声明式、高性能的方式完成任意时间网格上的累积统计——这才是真正的「云原生数据管道」写法。










