高效滚动非等值连接:大尺度金融数据匹配方案问询
超大规模金融数据:寻找每行历史最近更高ask值的rowid方案
核心思路:单调栈(O(n)时间复杂度)
这个问题本质是寻找前驱更大元素(Previous Greater Element, PGE),经典最优解法是单调栈,全程线性遍历,无需复杂join或嵌套查询,完美适配千万级数据量。
R 实现:data.table + Rcpp单调栈
利用data.table的高效分组能力,结合Rcpp实现的单调栈(C++级速度),解决单线程性能瓶颈:
library(data.table) library(Rcpp) # Rcpp实现单调栈,返回每行对应的历史最近更高ask的rowid cppFunction(' IntegerVector find_previous_greater(const NumericVector& ask, const IntegerVector& rowid) { int n = ask.size(); IntegerVector res(n, NA_INTEGER); std::stack<int> st; // 存储索引,栈内元素对应ask值单调递减 for (int i = 0; i < n; ++i) { // 弹出栈顶所有ask <= 当前ask的元素 while (!st.empty() && ask[st.top()] <= ask[i]) { st.pop(); } // 栈非空则栈顶就是目标索引 if (!st.empty()) { res[i] = rowid[st.top()]; } // 将当前索引压入栈 st.push(i); } return res; }') # 加载数据(假设数据包含rowid、ask、instrument_id列) data <- fread("your_data.csv") # 全局处理(无分组) data[, prev_greater_rowid := find_previous_greater(ask, rowid)] # 分组处理(按标的分组,比如instrument_id) data[, prev_greater_rowid := find_previous_greater(ask, rowid), by = .(instrument_id)] # 多线程优化:拆分分组并行计算(适合多分组场景) library(parallel) groups <- unique(data$instrument_id) result_list <- mclapply(groups, function(g) { sub_dt <- data[instrument_id == g] sub_dt[, prev_greater_rowid := find_previous_greater(ask, rowid)] sub_dt }, mc.cores = detectCores() - 1) data_parallel <- rbindlist(result_list)
Python 实现:Numba加速单调栈
避开Polars的join限制,用Numba将Python代码编译为机器码,速度接近C++,代码简洁易维护:
import numba import polars as pl import numpy as np @numba.jit(nopython=True) def find_previous_greater(ask, rowid): n = len(ask) res = np.full(n, -1, dtype=np.int64) # 用-1标记无匹配项 stack = [] for i in range(n): # 弹出栈顶所有小于等于当前ask的元素 while stack and ask[stack[-1]] <= ask[i]: stack.pop() if stack: res[i] = rowid[stack[-1]] stack.append(i) return res # 加载数据 df = pl.read_csv("your_data.csv") # 全局处理 ask_arr = df["ask"].to_numpy() rowid_arr = df["rowid"].to_numpy() df = df.with_columns(pl.Series("prev_greater_rowid", find_previous_greater(ask_arr, rowid_arr))) # 分组处理(按instrument_id) def process_group(group): ask = group["ask"].to_numpy() rowid = group["rowid"].to_numpy() return group.with_columns(pl.Series("prev_greater_rowid", find_previous_greater(ask, rowid))) df_grouped = df.group_by("instrument_id").apply(process_group)
C++/Rcpp 实现:原生单调栈(极致性能)
追求最快速度和最低内存占用,直接用C++实现逻辑,通过Rcpp嵌入R调用:
#include <Rcpp.h> #include <stack> using namespace Rcpp; // [[Rcpp::export]] IntegerVector previous_greater_rowid(NumericVector ask, IntegerVector rowid) { int n = ask.size(); IntegerVector res(n, NA_INTEGER); std::stack<int> idx_stack; for (int i = 0; i < n; ++i) { while (!idx_stack.empty() && ask[idx_stack.top()] <= ask[i]) { idx_stack.pop(); } if (!idx_stack.empty()) { res[i] = rowid[idx_stack.top()]; } idx_stack.push(i); } return res; }
编译后在R中直接调用,用法同第一个R方案。
方案对比
- R+Rcpp:兼顾R的数据分析生态和C++性能,适合熟悉R的用户,多线程分组可进一步提速。
- Python+Numba:无需C基础,代码简洁,编译后速度接近原生C,适配Python生态。
- 纯C++/Rcpp:速度最快、内存占用最低,适合极致性能需求场景。
内容的提问来源于stack exchange,提问作者gaut
相关产品推荐
相关产品推荐

