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

Polars中基于last_trade_ts触发的updated_at滚动分组问题

Python Polars 高效实现方案

核心思路:仅基于last_trade_ts的唯一值创建窗口,规避全量滚动的冗余计算

  1. 提取窗口触发行
    先筛选出last_trade_ts首次出现的行,这些行就是窗口的起始触发点:
import polars as pl

# 确保时间字段为datetime类型
df = df.with_columns(pl.col("updated_at").str.to_datetime())

# 获取每个last_trade_ts对应的首行
trigger_rows = df.unique(subset="last_trade_ts", keep="first")
  1. 生成窗口起止时间
    为每个触发行计算5分钟窗口的结束时间:
trigger_windows = trigger_rows.with_columns(
    window_end=pl.col("updated_at") + pl.duration(minutes=5)
).select("updated_at", "window_end", "last_trade_ts")
  1. 关联原数据到对应窗口
    通过范围连接将原数据分配到符合条件的窗口中(原数据的updated_at需落在窗口的起止区间内):
joined = df.join_asof(
    trigger_windows,
    left_on="updated_at",
    right_on="updated_at",
    strategy="backward",
    allow_parallel=True
).filter(pl.col("updated_at") <= pl.col("window_end"))
  1. 按窗口分组聚合
    基于窗口起始/结束时间执行自定义聚合逻辑:
result = joined.group_by("updated_at", "window_end").agg(
    count=pl.count(),
    avg_price=pl.col("trade_price").mean()
    # 替换为你的实际聚合需求
)

R 语言实现方案

用dplyr+fuzzyjoin或data.table实现高效非等连接

方法1:dplyr + fuzzyjoin

  1. 提取窗口触发行
library(dplyr)
library(fuzzyjoin)
library(lubridate)

# 转换时间字段格式
df <- df %>% mutate(updated_at = ymd_hms(updated_at))

# 获取每个last_trade_ts的首行
trigger_rows <- df %>% distinct(last_trade_ts, .keep_all = TRUE)
  1. 生成窗口并关联原数据
trigger_windows <- trigger_rows %>% 
    mutate(window_end = updated_at + minutes(5)) %>%
    select(updated_at, window_end, last_trade_ts)

# 模糊连接匹配窗口区间
joined <- fuzzy_left_join(
    df,
    trigger_windows,
    by = c("updated_at" = "updated_at", "updated_at" = "window_end"),
    match_fun = list(`>=`, `<=`)
)
  1. 分组聚合
result <- joined %>%
    group_by(window_start = updated_at.y, window_end) %>%
    summarise(
        count = n(),
        avg_price = mean(trade_price, na.rm = TRUE)
        # 替换为你的实际聚合需求
    ) %>%
    ungroup()

方法2:data.table 大数据场景更高效

library(data.table)
library(lubridate)

setDT(df)
df[, updated_at := ymd_hms(updated_at)]

# 提取触发行
trigger_rows <- unique(df, by = "last_trade_ts", fromLast = FALSE)
trigger_rows[, window_end := updated_at + minutes(5)]

# 非等连接匹配窗口
joined <- df[trigger_rows, 
             on = .(updated_at >= updated_at, updated_at <= window_end),
             allow.cartesian = TRUE]

# 分组聚合
result <- joined[, .(
    count = .N,
    avg_price = mean(trade_price, na.rm = TRUE)
), by = .(window_start = updated_at, window_end)]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 23:32:11