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

百万行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实现的分组操作,速度比原方案快数十倍
  • 内存占用更低,适合百万级甚至千万级数据集
  • 批量处理逻辑减少了重复的分组计算

总结

两种方案均通过避免逐行迭代解决了性能问题:

  1. 若偏好tidyverse语法,选择dplyr+slider方案,兼顾性能与可读性
  2. 若处理超大规模数据,优先选择data.table方案,性能优势更明显

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 13:20:31