【问题标题】:repeated joins on huge datasets在庞大的数据集上重复连接
【发布时间】:2016-04-19 05:59:37
【问题描述】:

几个2列的矩阵需要组合起来,如下图

matrix1
1,3
1,5
3,6

matrix2
1,4
1,5
3,6
3,7

output
1,3,1
1,4,1
1,5,2
3,6,2
3,7,1

输出中的第三列是对在所有矩阵中出现的次数的计数。我写了一些代码来做到这一点

require(data.table)

set.seed(1000)
data.lst <- lapply(1:200, function(n) { x <- matrix(sample(1:1000,2000,replace=T), ncol=2); x[!duplicated(x),] })

#method 1
pair1.dt <- data.table(i=integer(0), j=integer(0), cnt=integer(0))
for(mat in data.lst) {
    pair1.dt <- rbind(pair1.dt, data.table(i=mat[,1],j=mat[,2],cnt=1))[, .(cnt=sum(cnt)), .(i,j)]
}

#method 2
pair2.dt <- data.table(i=integer(0), j=integer(0), cnt=integer(0))
for(mat in data.lst) {
    pair2.dt <- merge(pair2.dt, data.table(i=mat[,1],j=mat[,2],cnt=1), by=c("i","j"), all=T)[, 
        cnt:=rowSums(.SD,na.rm=T), .SDcols=c("cnt.x","cnt.y")][, c("cnt.x","cnt.y"):=NULL]
}

cat(sprintf("num.rows  =>  pair1: %d,  pair2: %d", pair1.dt[,.N], pair2.dt[,.N]), "\n")

在真正的问题中,每个矩阵都有数以百万计的行,并且可能有 30-40% 的重叠。我试图找出最快的方法来做到这一点。我尝试使用 Matrix::sparseMatrix。虽然这要快得多,但我遇到了一个错误“尚不支持长向量”。我在这里有几种不同的基于 data.table 的方法。我正在寻找加快此代码速度和/或建议替代方法的建议。

【问题讨论】:

  • 一方面,不要创建cnt 并将其相加。相反,请尝试 rbindlist(lapply(data.lst, as.data.table))[, .N, by=V1:V2] 或类似名称。
  • 我认为这里不需要迭代,但对我来说你想要实现的目标并不明显,所以我可能错了。
  • 我在这里进行迭代的原因是,在 real 问题中,每个矩阵都有数千万行。所以,我无法保存所有这些矩阵。每次matrix 可用时,我都需要更新output
  • 我愿意。两者都是 1:200 万。数据非常非常稀疏。
  • @ironv 我有点困惑——如果你真的有这些矩阵的列表,那么你可以将所有这些存储在内存中,然后执行rbindlist正如弗兰克建议的那样……?

标签: r data.table


【解决方案1】:

首先,制作它们的data.tables:

dt.lst = lapply(data.lst, as.data.table)

堆叠。作为比较,这里是涉及堆叠的快速方法:

res0 = rbindlist(dt.lst)[, .(n = .N), by=V1:V2]

OP 说这是不可行的,因为rbindlist 的中间结果会太大。

先枚举。如果值范围很小,我建议先枚举它们:

res1 = CJ(V1 = 1:1000, V2 = 1:1000)[, n := 0L]
for (k in seq_along(dt.lst)) res1[ dt.lst[[k]], n := n + .N, by=.EACHI ] 

fsetequal(res0, res1[n>0]) # TRUE

OP 指出有 1e12 个可能的值,所以这似乎不是一个好主意。相反,我们可以使用

res2 = dt.lst[[1L]][0L]
for (k in seq_along(dt.lst)) res2 = funion(res2, dt.lst[[k]])
res2[, n := 0L]
setkey(res2, V1, V2)
for (k in seq_along(dt.lst)) res2[ dt.lst[[k]], n := n + .N, by=.EACHI ]     

fsetequal(res0, res2) # TRUE

对于给出的示例,这是三个选项中最慢的,但鉴于 OP 的担忧,这对我来说似乎是最好的。

在循环中增长。最后...

res3 = dt.lst[[1L]][0L][, n := NA_integer_][]
for (k in seq_along(dt.lst)){
  setkey(res3, V1, V2)
  res3[dt.lst[[k]], n := n + .N, by=.EACHI ]
  res3 = rbind(
    res3, 
    fsetdiff(dt.lst[[k]], res3[, !"n", with=FALSE], all=TRUE)[, .(n = .N), by=V1:V2]
  )
} 

fsetequal(res0, res3) # TRUE

在 R 中强烈建议不要在循环内增长对象且效率低下,但这允许您在一个循环中而不是两个循环中完成。

其他选项和注释。我怀疑您最好使用哈希。这些可以在 hash 包中获得,也可能通过 Rcpp 包获得。

fsetequalfsetdifffunion 是该软件包开发版的最新添加。在 data.table 项目的官方网站上查找详细信息。

顺便说一句,如果每个矩阵中的条目不同,您可以将上面的任何地方的.N 替换为1L 并删除by=.EACHIall=TRUE

【讨论】:

  • fsetequal 考虑到 nrow 已经在比较集合重复项。
  • @Frank 我实施了您的建议“在循环中增长”以及由于矩阵条目是唯一的更改(您帖子的最后一行)。在真正的问题中,到第 20 次迭代(即使用 dt.lst 中的第 20 个条目),setkey 需要 159 秒,n 递增需要 6 秒,fsetdiff 需要 240 秒,rbind 需要 80 秒。我想知道您关于哈希的建议,它会与 data.table 一起使用吗?
  • @ironv 不一定需要涉及 data.table。我认为斯科特的回答符合我的想法。 (我对 Rcpp 不够熟悉,不能肯定地说。)如果您尝试了双循环方式(在此答案中使用 res2),我想知道这是怎么回事。
【解决方案2】:

使用 Rcpp。此方法将利用 std::unordered_map 的散列属性。

#include "Rcpp.h"
#include <stdint.h>
#include <unordered_map>

using namespace std;
using namespace Rcpp;

//[[Rcpp::export]]
Rcpp::XPtr<int> CreateMap(){
  std::unordered_map<int64_t, int>* myMap = new std::unordered_map<int64_t, int>();
  Rcpp::XPtr<int> p((int*)myMap,false);
  return p;
}

//[[Rcpp::export]]
void FreeMap(Rcpp::XPtr<int> map_ptr){
  std::unordered_map<int64_t, int>* myMap =  (std::unordered_map<int64_t, int>*)(int*)map_ptr;
  delete myMap;
}

//[[Rcpp::export]]
void AccumulateValues(Rcpp::XPtr<int> map_ptr, SEXP mat){

  NumericMatrix m(mat);

  std::unordered_map<int64_t, int>* myMap =  (std::unordered_map<int64_t, int>*)(int*)map_ptr;
  for(int i = 0; i<m.nrow(); i++){
    int c1 = m(i, 0);
    int c2 = m(i, 1);
    int64_t key = ((int64_t)c1 << 32) + c2;
    (*myMap)[key] ++;

  }
}
//[[Rcpp::export]]
SEXP AsMatrix(Rcpp::XPtr<int> map_ptr){
  std::unordered_map<int64_t, int>* myMap =  (std::unordered_map<int64_t, int>*)(int*)map_ptr;
  NumericMatrix m(myMap->size(),3);
  int index = 0;
  for ( auto it = myMap->begin(); it != myMap->end(); ++it ){
    int64_t key = it->first;
    m(index, 0) = (int)(key >> 32);
    m(index, 1) = (int)key;
    m(index, 2) = it->second;
    index++;
  }
  return m;
}

R 代码是:

myMap<-CreateMap()
AccumulateValues(myMap, matrix1)
AccumulateValues(myMap, matrix2)
result<-AsMatrix(myMap)
FreeMap(myMap)

也需要

PKG_CXXFLAGS = "-std=c++0x"

在 makevars 包中

【讨论】:

  • 感谢您的帖子。我从未使用过 Rcpp 并且真的不知道如何使用上面的代码。将查看示例并尝试一下。
  • 我要做的步骤(可能有很多方法)将是 1. 确保您已正确安装和配置 R 工具 2. 安装 Rcpp 包 3. 通过 R 创建一个新项目工作室在对话框中指定新的“带有 Rcpp 的项目” 4. 在项目目录中,您现在有一个带有“hello world”示例的“src”目录,将您要使用的 C++ 源代码放在那里 5.因为上面的代码使用unordered_map 您需要使用编译器标志来支持它,因此将名为 Makevars 的文件添加到“src”并添加 PKG_CXXFLAGS = "-std=c++0x" 行。
  • 感谢您的帖子。我刚刚开始使用 Rcpp 并尝试了 Dirk 网站上的斐波那契示例。惊人的加速!我在没有 root 权限的 Linux 机器上运行它。我已经编译了 R 并安装了 Rcpp。我是在编译包时添加 Makevars 文件还是只为上面的代码添加?
  • 由于我概述的步骤是基于包的方法(即,您正在创建和编译您自己的专用于此代码的“临时”R 包),您只需要确保当你编译你的包(或任何使用更新的 C++ 特性的包)时,Makevars 文件存在于你包的 src 目录中如果你使用内联 Rcpp 特性而不是包方法,请参阅这篇文章stackoverflow.com/questions/7063265/…跨度>
  • 在 R 控制台中,我使用 Sys.setenv("PKG_CXXFLAGS"="-std=c++0x")sourceCpp 编译和使用代码。每个矩阵有 8000 万行,Accumulate 花费了大约 25 秒。谢谢您的帮助。你的回答不仅帮我解决了这个问题,还帮我学习了一个非常有用的包!
【解决方案3】:

也许您可以在内存允许的情况下分批处理数据:

maxRows = 5000 # approximately how many rows can you hold in memory
tmp.lst = list()
nrows = 0
idx = 1
for (i in seq_along(data.lst)) {
  tmp.lst[[idx]] = as.data.table(data.lst[[i]])[, cnt := 1]
  idx = idx + 1
  nrows = nrows + nrow(data.lst[[i]])

  # if too many rows, collapse (can also replace by some memory condition)
  if (nrows > maxRows) {
    tmp.lst = list(rbindlist(tmp.lst)[, .(cnt = sum(cnt)), by = V1:V2])
    idx = 2
    nrows = nrow(tmp.lst[[1]])
  }
}

#final collapse
res = rbindlist(tmp.lst)[, .(cnt = sum(cnt)), by = V1:V2]

【讨论】:

    猜你喜欢
    • 2012-09-23
    • 1970-01-01
    • 2023-04-04
    • 1970-01-01
    • 1970-01-01
    • 2011-04-23
    • 1970-01-01
    • 2016-04-19
    • 2022-01-18
    相关资源
    最近更新 更多