Polars中基于last_trade_ts触发的updated_at滚动分组问题
Python Polars 高效实现方案
核心思路:仅基于last_trade_ts的唯一值创建窗口,规避全量滚动的冗余计算
- 提取窗口触发行
先筛选出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")
- 生成窗口起止时间
为每个触发行计算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")
- 关联原数据到对应窗口
通过范围连接将原数据分配到符合条件的窗口中(原数据的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"))
- 按窗口分组聚合
基于窗口起始/结束时间执行自定义聚合逻辑:
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
- 提取窗口触发行
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)
- 生成窗口并关联原数据
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(`>=`, `<=`) )
- 分组聚合
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
相关产品推荐
相关产品推荐

