【问题标题】:Spark out of memory when reducing by key按键减少时内存不足
【发布时间】:2018-02-20 08:47:14
【问题描述】:

我正在研究一种需要对大型矩阵进行数学运算的算法。基本上,该算法涉及以下步骤:

输入:大小为 n 的两个向量 u 和 v

  1. 对于每个向量,计算向量中元素之间的成对欧几里得距离。返回两个矩阵 E_u 和 E_v

  2. 对于两个矩阵中的每个条目,应用一个函数 f。返回两个矩阵M_u、M_v

  3. 求 M_u 的特征值和特征向量。返回 e_i, ev_i for i = 0,...,n-1

  4. 计算每个特征向量的外积。返回一个矩阵 O_i = e_i*transpose(e_i), i = 0,...,n-1

  5. 用 e_i = e_i + delta_i 调整每个特征值,其中 delta_i = 所有元素之和(O_i 和 M_v 的元素乘积)/2*mu,其中 mu 是一个参数

  6. 最终返回一个矩阵 A = elementwise sum (e_i * O_i) over i = 0,...,n-1

我面临的问题主要是当 n 很大(15000 或更大)时的内存,因为这里的所有矩阵都是密集矩阵。我目前的实现方式可能不是最好的,并且部分有效。

我对 M_u 使用了 RowMatrix 并使用 SVD 进行特征分解。

SVD 得到的 U 因子是一个行矩阵,其列是 ev_i 的,所以我必须手动转置它,使其行变为 ev_i。得到的 e 向量就是特征值 e_i。

由于之前尝试将每一行 ev_i 直接映射到 O_i 的尝试由于内存不足而失败,我目前正在做

R = U.map{
    case(i,ev_i) => {
      (i, ev_i.toArray.zipWithIndex)
    }
  }//add index for each element in a vector
  .flatMapValues(x=>x)}
  .join(U)//eigen vectors column is appended
  .map{case(eigenVecId, ((vecElement,elementId), eigenVec))=>(elementId, (eigenVecId, vecElement*eigenVec))}

为了在上面的步骤 5 中计算调整后的 e_i,M_v 存储为元组的 rdd (i,denseVector)。那么

deltaRdd = R.join(M_v)
  .map{
    case(j,((i,row_j_of_O_i),row_j_of_M_v))=>
    (i,row_j_of_O_i.t*DenseVector(row_j_of_M_v.toArray)/(2*mu))
  }.reduceByKey(_+_)

最后,为了计算 A,同样由于内存问题,我必须首先连接来自不同 rdds 的行,然后按键减少。具体来说,

R_rearranged = R.map{case(j, (i, row_j_of_O_i))=>(i,(j,row_j_of_O_i))}
termsForA = R_rearranged.join(deltaRdd)
A = termsForA.map{
  case(i,(j,row_j_of_O_i), delta_i)) => (j, (delta_i + e(i))*row_j_of_O_i)
}
.reduceByKey(_+_)

上面的实现工作到了termsForA的步骤,这意味着如果我对termsForA执行一个动作,比如termsForA.take(1).foreach(println),它就会成功。但是如果我在 A 上执行一个动作,比如 A.count(),就会在驱动程序上发生 OOM 错误。

我尝试调整 sparks 配置以增加驱动程序内存和并行度,但都失败了。

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    使用 IndexedRowMatrix 代替 RowMatrix,它将有助于转换和转置。 假设您的 IndexedRowMatrix 是 Irm

    svd = Irm.computeSVD(k, True)
    U = svd.U
    U =  U.toCoordinateMatrix().transpose().toIndexedRowMatrix()
    

    您可以将 Irm 转换为 BlockMatrix,以便与另一个分布式 BlockMatrix 相乘。

    【讨论】:

      【解决方案2】:

      我猜在某个时候,Spark 决定不需要对 executor 执行操作,而是在 driver 上完成所有工作。实际上,termForA 在计数等操作中也会失败。不知何故,我通过广播 deltaRdd 和 e 使它工作。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2013-11-19
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2014-05-03
        • 2015-07-08
        • 1970-01-01
        • 2013-05-21
        相关资源
        最近更新 更多