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

高效滚动非等值连接:大尺度金融数据匹配方案问询

超大规模金融数据:寻找每行历史最近更高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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 10:53:11