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

高效合并两个DataFrame:时间窗口与事件数据关联优化

优化时间窗口与事件数据的聚合方案(避免内存溢出)

我需要将两个数据集合并:一个是包含起止时间的时间窗口数据集(100-500个窗口),另一个是带时间戳的事件数据集(1000万条),最终生成每个窗口内各类事件的计数表。原方案使用cross join导致内存爆炸,进程被终止,需要更高效的实现方式。

原问题代码

import polars as pl

actions = {
    "id": ["a", "a", "a", "a", "b", "b", "a", "a"],
    "action": ["start", "stop", "start", "stop", "start", "stop", "start", "stop"],
    "time": [0.0, 1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0],
}

events = {
    "name": ["x", "x", "x", "y", "y", "z", "w", "w", "w"],
    "time": [0.0, 0.1, 0.5, 1.1, 2.5, 3.0, 4.5, 4.9, 5.5],
}

actions_df = (
    pl.DataFrame(actions)
    .group_by("id")
    .agg(
        start=pl.col("time").filter(pl.col("action") == "start"),
        stop=pl.col("time").filter(pl.col("action") == "stop"),
    )
    .explode(["start", "stop"])
)

df = (
    actions_df.join(pl.DataFrame(events), how="cross")
    .filter((pl.col("time") >= pl.col("start")) & (pl.col("time") <= pl.col("stop")))
    .group_by(["id", "start", "stop", "name"])
    .agg(count=pl.count("name"))
    .pivot("name", index=["id", "start", "stop"], values="count")
    .fill_null(0)
)

result_df = (
    actions_df.join(df, on=["id", "start", "stop"], how="left")
    .fill_null(0)
    .sort("start")
)

print(result_df)

优化思路

避免生成笛卡尔积,利用事件数据按时间排序后的二分查找特性,快速定位每个窗口对应的事件区间,再统计区间内的事件类型计数。该方法内存占用极低,仅需处理事件数据一次,每个窗口通过两次二分查找定位事件范围,时间复杂度为O(N log N + W)(N为事件数,W为窗口数)。

优化代码

import polars as pl
import numpy as np

# 1. 预处理窗口数据:生成id+起止时间的窗口表
actions = {
    "id": ["a", "a", "a", "a", "b", "b", "a", "a"],
    "action": ["start", "stop", "start", "stop", "start", "stop", "start", "stop"],
    "time": [0.0, 1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0],
}

events = {
    "name": ["x", "x", "x", "y", "y", "z", "w", "w", "w"],
    "time": [0.0, 0.1, 0.5, 1.1, 2.5, 3.0, 4.5, 4.9, 5.5],
}

actions_df = (
    pl.DataFrame(actions)
    .group_by("id")
    .agg(
        start=pl.col("time").filter(pl.col("action") == "start"),
        stop=pl.col("time").filter(pl.col("action") == "stop"),
    )
    .explode(["start", "stop"])
    .with_row_index("window_idx")  # 添加窗口索引,用于后续合并
)

# 2. 预处理事件数据:按时间排序,提取numpy数组加速查询
events_df = pl.DataFrame(events).sort("time")
event_times = events_df["time"].to_numpy()
event_names = events_df["name"].to_numpy()

# 3. 对每个窗口统计事件数量
def count_events_in_window(row):
    start = row["start"]
    stop = row["stop"]
    # 二分查找定位事件区间
    left_idx = np.searchsorted(event_times, start, side="left")
    right_idx = np.searchsorted(event_times, stop, side="right")
    # 提取区间内的事件名称并统计计数
    if left_idx >= right_idx:
        return pl.DataFrame({"window_idx": [row["window_idx"]]})
    name_counts = pl.Series(event_names[left_idx:right_idx]).value_counts()
    return pl.DataFrame({"window_idx": [row["window_idx"]]}).hstack(name_counts)

# 遍历所有窗口,合并统计结果
window_counts = actions_df.map_rows(count_events_in_window).group_by("window_idx").first()

# 4. 合并窗口数据与统计结果,补全缺失事件类型的0值
result_df = (
    actions_df.join(window_counts, on="window_idx", how="left")
    .drop("window_idx")
    .fill_null(0)
    .sort("start")
)

print(result_df)

代码说明

  • 窗口预处理:生成标准的窗口表并添加索引,方便后续合并。
  • 事件预处理:按时间排序后转为numpy数组,利用numpy的二分查找searchsorted快速定位窗口对应的事件区间,比纯Polars操作更快。
  • 窗口事件统计:对每个窗口,通过二分查找获取事件区间,统计区间内各事件类型的数量,无事件时返回空表。
  • 结果合并:将统计结果与原窗口表左连接,补全缺失事件类型的0值,得到最终结果。

优势

  • 内存占用极低:无需生成笛卡尔积,仅需保存事件数据和窗口数据的原始大小。
  • 速度高效:二分查找定位事件区间,避免全量数据匹配,适合千万级事件数据。
  • 支持重叠窗口:一个事件可被多个包含它的窗口统计,符合实际业务需求。

内容的提问来源于stack exchange,提问作者DJDuque

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 01:17:02