在Polars中实现非等值连接:统计时间窗口内事件数量
使用Polars统计时间窗口内的事件频次
我需要用Polars处理一个事件统计问题:给定事件列表(含名称和时间戳)和时间窗口列表,统计每个窗口内各类事件的发生次数。目前已经通过join_asof拿到了每个窗口对应的事件首尾索引,接下来需要完成频次统计。
现有代码
import polars as pl events = { "name": ["a", "b", "a", "b", "a", "c", "b", "a", "b", "a", "b", "a", "b", "a", "b", "a", "b", "a", "b"], "time": [0.0, 1.0, 1.5, 2.0, 2.25, 2.26, 2.45, 2.5, 3.0, 3.4, 3.5, 3.6, 3.65, 3.7, 3.8, 4.0, 4.5, 5.0, 6.0], } windows = { "start_time": [1.0, 2.0, 3.0, 4.0], "stop_time": [3.5, 2.5, 3.7, 5.0], } events_df = pl.DataFrame(events).sort("time").with_row_index() windows_df = ( pl.DataFrame(windows) .sort("start_time") .join_asof(events_df, left_on="start_time", right_on="time", strategy="forward") .drop("name", "time") .rename({"index": "first_index"}) .sort("stop_time") .join_asof(events_df, left_on="stop_time", right_on="time", strategy="backward") .drop("name", "time") .rename({"index": "last_index"}) ) print(windows_df)
输出结果:
shape: (4, 4) ┌────────────┬───────────┬─────────────┬────────────┐ │ start_time ┆ stop_time ┆ first_index ┆ last_index │ │ --- ┆ --- ┆ --- ┆ --- │ │ f64 ┆ f64 ┆ u32 ┆ u32 │ ╞════════════╪═══════════╪═════════════╪════════════╡ │ 2.0 ┆ 2.5 ┆ 3 ┆ 7 │ │ 1.0 ┆ 3.5 ┆ 1 ┆ 10 │ │ 3.0 ┆ 3.7 ┆ 8 ┆ 13 │ │ 4.0 ┆ 5.0 ┆ 15 ┆ 17 │ └────────────┴───────────┴─────────────┴────────────┘
预期输出
shape: (4, 5) ┌────────────┬───────────┬─────┬─────┬─────┐ │ start_time ┆ stop_time ┆ a ┆ b ┆ c │ │ --- ┆ --- ┆ --- ┆ --- ┆ --- │ │ f64 ┆ f64 ┆ i64 ┆ i64 ┆ i64 │ ╞════════════╪═══════════╪═════╪═════╪═════╡ │ 1.0 ┆ 3.5 ┆ 4 ┆ 5 ┆ 1 │ │ 2.0 ┆ 2.5 ┆ 2 ┆ 2 ┆ 1 │ │ 3.0 ┆ 3.7 ┆ 3 ┆ 3 ┆ 0 │ │ 4.0 ┆ 5.0 ┆ 2 ┆ 1 ┆ 0 │ └────────────┴───────────┴─────┴─────┴─────┘
解决方案
可以通过int_ranges生成每个窗口覆盖的事件索引范围,结合explode展开后关联事件数据,最后分组透视得到结果:
# 生成每个窗口对应的索引范围,展开后关联事件名称 result = ( windows_df .with_columns(pl.int_ranges("first_index", pl.col("last_index") + 1).alias("event_indices")) .explode("event_indices") .join(events_df.select("index", "name"), left_on="event_indices", right_on="index") # 按窗口分组,统计各事件数量 .group_by("start_time", "stop_time") .agg( pl.col("name").value_counts(sort=False).alias("counts") ) # 透视成宽表,填充缺失事件的计数为0 .unnest("counts") .pivot( index=["start_time", "stop_time"], columns="name", values="count", aggregate_function="sum" ) .fill_null(0) # 按start_time排序匹配预期输出 .sort("start_time") ) print(result)
代码说明
- 生成索引范围:用
int_ranges生成从first_index到last_index(含)的索引序列,注意要给last_index加1,因为int_ranges是左闭右开区间。 - 展开索引:用
explode将每个窗口的索引序列展开成多行,每行对应一个事件索引。 - 关联事件数据:通过索引关联事件表,获取每个索引对应的事件名称。
- 分组统计:按窗口的
start_time和stop_time分组,用value_counts统计每组内各事件的出现次数。 - 透视表转换:将长表转换为宽表,把事件名称作为列,计数作为值,缺失的事件用
fill_null(0)补0。 - 排序:最后按
start_time排序,和预期输出格式一致。
内容的提问来源于stack exchange,提问作者DJDuque
相关产品推荐
相关产品推荐

