
本文介绍如何使用 Polars 原生表达式(无需 map_elements 或 Python 回调)高效计算事务流在自定义时间网格上的运行累计和,核心是 join_asof + 左连接 + 分组聚合 + cum_sum。
本文介绍如何使用 polars 原生表达式(无需 `map_elements` 或 python 回调)高效计算事务流在自定义时间网格上的运行累计和,核心是 `join_asof` + 左连接 + 分组聚合 + `cum_sum`。
在时间序列分析中,常需将离散事件(如交易流水)映射到规则时间网格(如每分钟一个时间点),并计算截至每个网格点的累计值(例如“截至 11:03 的总金额”)。若采用 map_elements 配合 filter 手动遍历(如原示例中的 _cumsum 函数),不仅性能差(Python 层循环 + 多次 DataFrame 过滤),还无法利用 Polars 的惰性执行与并行优化。
更优解是完全基于 Polars 原生表达式链式操作,分四步完成:
对齐事件到目标时间网格:使用
join_asof(..., strategy="forward"),为每个原始事件(df)匹配dg中首个不早于该事件时间的时间点。这避免了手动循环,且底层由 Rust 高效实现。反向关联金额至网格:将上一步结果重命名
time_right → time后,与dg左连接(how="left"),使每个网格时间点携带其“应计入”的所有事件金额(未匹配则为null)。处理重复与缺失:因多个事件可能落入同一网格点(如两笔交易均发生在 11:02:30,均映射至
11:03:00),需按time分组并求和;同时,left join产生的null值需填充为0.0(可显式用.fill_null(0)或依赖sum()的默认行为——空组返回0)。计算累积和:最后对聚合后的
amount列调用.cum_sum(),并重命名为cum_amount。
完整代码如下(含关键注释):
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(
pl.datetime_range(
datetime.datetime(2025, 2, 2, 11, 0),
datetime.datetime(2025, 2, 2, 11, 5),
"1m",
eager=True
),
schema=["time"]
)
# 四步链式表达式:高效、可读、可优化
result = (
dg
.join(
df.join_asof(
dg,
on="time",
strategy="forward", # 关键:向前查找最近网格点
coalesce=False
)
.select("amount", pl.col("time_right").alias("time")), # 重命名对齐列
on="time",
how="left"
)
.group_by("time")
.agg(pl.col("amount").sum().fill_null(0)) # 显式填充 null 为 0,增强鲁棒性
.sort("time") # 确保时间顺序(虽通常已有序,但 group_by 后建议显式排序)
.with_columns(
pl.col("amount").cum_sum().alias("cum_amount")
)
)
print(result)
输出:
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 │ └─────────────────────┴────────┴────────────┘
✅ 优势总结:
-
高性能:全程使用 Polars 原生操作,避免 Python 解释器开销;
join_asof经过高度优化,复杂度接近 O(n+m); -
可扩展:天然支持大数据集(结合
lazy()惰性模式); - 健壮性:自动处理边界情况(如无事件的起始点、多事件同网格、时间超出范围等);
- 清晰语义:每步操作意图明确,易于调试与复用。
⚠️ 注意事项:
-
join_asof要求on列已排序(df和dg的time列需升序);若未排序,先调用.sort("time"); - 若需“截至但不包含”某时间点(如“11:03 的累计值仅含 11:03 之前的数据”),应将
strategy改为"backward"并调整时间偏移,或在join_asof前对dg时间减去微小量; - 生产环境中建议对
amount列添加.cast(pl.Float64)确保数值类型一致,避免隐式转换异常。










