如何基于时间差移除Polars数据框中的重复产品浏览记录?
移除产品浏览记录中10分钟内的重复条目(保留元数据)
需求说明
处理带时间戳的产品浏览数据时,需要移除同产品下与上一条非重复记录间隔10分钟内的重复条目(比如用户刷新页面产生的重复记录),同时保留每条记录的metadata。具体规则:
- 同产品首次浏览后,10分钟内的再次浏览视为重复,需移除
- 间隔超过10分钟的浏览需保留
- 不同产品的浏览记录互不影响
示例数据与期望输出
原始数据
from datetime import datetime import polars as pl df = pl.DataFrame( { "created_time": [ datetime(2023, 1, 1, 0, 0), datetime(2023, 1, 1, 0, 1), datetime(2023, 1, 1, 0, 2), datetime(2023, 1, 1, 0, 3), datetime(2023, 1, 1, 0, 11), datetime(2023, 1, 1, 0, 29), datetime(2023, 1, 1, 0, 31), ], "product_id": [1, 1, 2, 1, 1, 1, 1], "metadata":["a", "b", "c", "d", "e", "f", "g"] } )
打印结果:
shape: (7, 3) ┌─────────────────────┬────────────┬──────────┐ │ created_time ┆ product_id ┆ metadata │ │ --- ┆ --- ┆ --- │ │ datetime[μs] ┆ i64 ┆ str │ ╞═════════════════════╪════════════╪══════════╡ │ 2023-01-01 00:00:00 ┆ 1 ┆ a │ │ 2023-01-01 00:01:00 ┆ 1 ┆ b │ │ 2023-01-01 00:02:00 ┆ 2 ┆ c │ │ 2023-01-01 00:03:00 ┆ 1 ┆ d │ │ 2023-01-01 00:11:00 ┆ 1 ┆ e │ │ 2023-01-01 00:29:00 ┆ 1 ┆ f │ │ 2023-01-01 00:31:00 ┆ 1 ┆ g │ └─────────────────────┴────────────┴──────────┘
期望输出
df_desirable = pl.DataFrame( { "created_time": [ datetime(2023, 1, 1, 0, 0), datetime(2023, 1, 1, 0, 2), datetime(2023, 1, 1, 0, 11), datetime(2023, 1, 1, 0, 29), ], "product_id": [1, 2, 1, 1], "metadata":["a", "c", "e", "f"] } )
打印结果:
shape: (4, 3) ┌─────────────────────┬────────────┬──────────┐ │ created_time ┆ product_id ┆ metadata │ │ --- ┆ --- ┆ --- │ │ datetime[μs] ┆ i64 ┆ str │ ╞═════════════════════╪════════════╪══════════╡ │ 2023-01-01 00:00:00 ┆ 1 ┆ a │ │ 2023-01-01 00:02:00 ┆ 2 ┆ c │ │ 2023-01-01 00:11:00 ┆ 1 ┆ e │ │ 2023-01-01 00:29:00 ┆ 1 ┆ f │ └─────────────────────┴────────────┴──────────┘
问题分析
之前尝试的group_by("product_id")或group_by_dynamic无法满足需求:
- 简单分组后用
shift()只能对比同组的上一条记录,而非上一条被保留的记录,会误判间隔 group_by_dynamic是按固定时间窗口聚合,无法灵活追踪每条记录与上一条有效记录的间隔
解决方案
方法1:分组后自定义过滤函数(直观易懂)
通过对每个产品分组,逐行判断是否与上一条保留记录间隔超过10分钟,符合条件则保留:
from datetime import datetime, timedelta import polars as pl # 确保数据按时间排序 df = df.sort("created_time") def filter_recent_views(group): if len(group) == 0: return group # 初始化保留第一条记录 keep_indices = [0] last_kept_time = group[0, "created_time"] for idx in range(1, len(group)): time_diff = group[idx, "created_time"] - last_kept_time # 间隔超过10分钟则保留,并更新上一条保留时间 if time_diff > timedelta(minutes=10): keep_indices.append(idx) last_kept_time = group[idx, "created_time"] return group[keep_indices] # 分组应用过滤函数,保持原始顺序 result = df.group_by("product_id", maintain_order=True).apply(filter_recent_views) # 最后按时间重新排序(分组后可能打乱顺序) result = result.sort("created_time") print(result)
方法2:向量化窗口函数(高效处理大数据)
利用Polars的窗口函数和状态变量,实现更高效的向量化处理:
from datetime import timedelta import polars as pl df = df.sort("created_time") result = ( df .with_row_index() .group_by("product_id", maintain_order=True) .agg( pl.all(), # 计算每条记录与上一条保留记录的时间差 pl.col("created_time") .cumulative_eval( lambda s: pl.when(s.len() == 1) .then(pl.duration(seconds=0)) .otherwise(s.last() - pl.col("last_kept").last()), state=pl.col("created_time").alias("last_kept") ) .alias("time_since_last_kept") ) .explode(pl.exclude("time_since_last_kept")) .with_columns( # 标记需要保留的记录:第一条或间隔超10分钟 pl.when(pl.col("time_since_last_kept") > timedelta(minutes=10)) .or(pl.col("time_since_last_kept") == timedelta(seconds=0)) .then(True) .otherwise(False) .alias("keep") ) .filter(pl.col("keep")) .drop("index", "time_since_last_kept", "keep") .sort("created_time") ) print(result)
两种方法都能得到符合期望的输出,方法1适合小数据集,代码直观;方法2适合大数据集,性能更优。
内容的提问来源于stack exchange,提问作者TomNorway
相关产品推荐
相关产品推荐

