百万行R数据框时间序列特征聚合的高效实现方案问询
百万行DataFrame的高效时间序列聚合优化方案
问题背景
我有一个近百万行的DataFrame,包含a、b两类观测值,由两类测量数据关联生成,is_a、is_b列标记数据来源。需求如下:
- 按
ID分组 - 每个时间步聚合最近n个a、b值,不足n个时填充NA
- 同时计算每个值与当前时间步的小时差
tdiff
原方案使用purrr逐行迭代实现,处理大数据时速度极慢,原代码及示例数据如下:
示例数据
df <- tibble(ID = c(1, 1, 1, 1, 1), Time = c(as.POSIXct("1900-01-01 10:00:00"), as.POSIXct("1900-01-01 14:00:00"), as.POSIXct("1900-01-01 15:00:00"), as.POSIXct("1900-01-01 17:00:00"), as.POSIXct("1900-01-01 20:00:00")), a = c(NA, 1, 2, NA, 3), is_a = c(0, 1, 1, 0, 1), b = c(1, NA, NA, 3, 1), is_b = c(1, 0, 0, 1, 1))
原实现代码
aggregate_df <- function(use_data) { ret <- use_data %>% group_by(ID) %>% mutate(obs = purrr::pmap(.l = list(Time, ID), .f = function(curr_t, id) { df1 <- compute_lagged_sequence(., id, curr_t, "a", "is_a") df2 <- compute_lagged_sequence(., id, curr_t, "b", "is_b") return(tibble(df1, df2)) })) %>% unnest(obs) %>% select(-a, -is_a, -b, -is_b) return(ret) } compute_lagged_sequence <- function(use_data, id, curr_t, col_to_use, reference_col, n = 3) { rel_rows <- filter(use_data, ID == id) %>% filter(Time <= curr_t) %>% filter(.data[[reference_col]] == 1) %>% slice_tail(n = n) %>% ungroup() %>% select(Time, !!col_to_use) %>% mutate(!!sym(sprintf("%s_tdiff", col_to_use)) := as.double.difftime(curr_t - Time, units = "hours")) %>% select(-Time) %>% summarise(cur_data()[seq(n), ]) %>% arrange(!is.na(.data[[col_to_use]]), desc(!!sym(sprintf("%s_tdiff", col_to_use)))) %>% mutate(row = row_number()) %>% pivot_wider(names_from = row, values_from = c(!!col_to_use, !!sym(sprintf("%s_tdiff", col_to_use)))) return(rel_rows) }
优化思路
原方案的核心性能瓶颈是逐行调用函数执行过滤、切片、重组操作,这种循环在百万行数据下会产生大量重复计算。优化方向是利用向量化操作和高效分组工具,避免逐行迭代,具体可选择以下两种方案:
方案1:dplyr + Slider 向量化实现
使用slider包的滑动窗口函数替代逐行迭代,结合dplyr的分组能力,大幅减少重复计算:
library(dplyr) library(slider) library(tidyr) n <- 3 # 预处理:提前提取a、b的有效观测并排序 a_valid <- df %>% filter(is_a == 1) %>% select(ID, Time, a) %>% arrange(ID, Time) b_valid <- df %>% filter(is_b == 1) %>% select(ID, Time, b) %>% arrange(ID, Time) # 定义特征处理函数 process_feature <- function(main_df, valid_data, feat_name, n) { main_df %>% group_by(ID) %>% mutate( # 滑动窗口获取当前时间及之前的最近n个有效观测 feat_window = slide( .x = list(curr_time = Time), .f = function(curr) { valid_data %>% filter(ID == curr$ID[[1]], Time <= curr$curr_time[[1]]) %>% slice_tail(n = n) %>% mutate( tdiff = as.double.difftime(curr$curr_time[[1]] - Time, units = "hours"), row = row_number() ) %>% # 补全到n行,不足的填充NA complete(row = 1:n) %>% arrange(row) %>% select(!!feat_name, tdiff) }, .before = Inf ) ) %>% unnest(feat_window) %>% pivot_wider( names_from = row, values_from = c(!!feat_name, tdiff), names_glue = "{feat_name}_{.value}_{row}" ) %>% ungroup() } # 合并a、b特征结果 final_result <- df %>% select(ID, Time) %>% process_feature(a_valid, "a", n) %>% left_join(process_feature(df %>% select(ID, Time), b_valid, "b", n), by = c("ID", "Time")) %>% # 调整列顺序与原输出一致 select( ID, Time, a_1, a_2, a_3, a_tdiff_1, a_tdiff_2, a_tdiff_3, b_1, b_2, b_3, b_tdiff_1, b_tdiff_2, b_tdiff_3 ) print(final_result)
优势
- 滑动窗口操作是向量化的,避免逐行循环
- 提前预处理有效数据,减少重复过滤开销
- 代码风格与原方案接近,易理解和维护
方案2:data.table 高性能实现
data.table的分组和窗口操作基于C语言实现,性能远优于dplyr的逐行迭代,适合超大数据集:
library(data.table) n <- 3 setDT(df) # 预处理有效数据 a_dt <- df[is_a == 1, .(ID, Time, a)][order(ID, Time)] b_dt <- df[is_b == 1, .(ID, Time, b)][order(ID, Time)] # 定义处理函数 process_dt <- function(main_dt, feat_dt, feat_name, n) { main_dt[, { curr_id <- .BY$ID curr_times <- Time # 批量处理每个时间点的最近n个观测 lapply(curr_times, function(curr_t) { # 获取当前ID、当前时间之前的最近n个有效观测 rows <- feat_dt[ID == curr_id & Time <= curr_t][order(-Time)][1:n] # 补全n行 if (nrow(rows) < n) { rows <- rbind(rows, as.data.table(matrix(NA, ncol = 3, nrow = n - nrow(rows)))) setnames(rows, names(feat_dt)) } # 计算时间差并按时间升序排列(对应原输出的a_1到a_3) rows[, tdiff := as.double.difftime(curr_t - Time, units = "hours")] rows <- rows[order(Time)] # 整理为单行 data.table(t(c(unlist(rows[, get(feat_name)]), unlist(rows[, tdiff])))) }) %>% rbindlist() }, by = ID] %>% cbind(main_dt[, .(Time)]) %>% setnames(c("ID", sprintf("%s_%d", feat_name, 1:n), sprintf("%s_tdiff_%d", feat_name, 1:n), "Time")) %>% setcolorder(c("ID", "Time", sprintf("%s_%d", feat_name, 1:n), sprintf("%s_tdiff_%d", feat_name, 1:n))) } # 合并结果 result_dt <- process_dt(df[, .(ID, Time)], a_dt, "a", n) %>% merge(process_dt(df[, .(ID, Time)], b_dt, "b", n), by = c("ID", "Time")) print(result_dt)
优势
- 底层C实现的分组操作,速度比原方案快数十倍
- 内存占用更低,适合百万级甚至千万级数据集
- 批量处理逻辑减少了重复的分组计算
总结
两种方案均通过避免逐行迭代解决了性能问题:
- 若偏好tidyverse语法,选择dplyr+slider方案,兼顾性能与可读性
- 若处理超大规模数据,优先选择data.table方案,性能优势更明显
内容的提问来源于stack exchange,提问作者Ai4l2s
相关产品推荐
相关产品推荐

