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

如何基于时间差移除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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 21:25:11