【问题标题】:Store Spark distributed matrix in MongoDB在 MongoDB 中存储 Spark 分布式矩阵
【发布时间】:2016-09-13 21:54:52
【问题描述】:

在计算与存储在 HDFS 文件中的一组点相关的距离矩阵后,我需要通过MongoDB Connector for Apache Spark 将计算出的距离矩阵以分布式形式(CoordinateMatrix/RowMatrix)存储在 MongoDB 中。是否有推荐的方法来执行此操作,甚至有更好的连接器来进行此类操作?

这是我的代码的一部分:

val data = sc.textFile("hdfs://localhost:54310/usrp/copy_sample_data.txt")
val points = data.map(s => Vectors.dense(s.split(',').map(_.toDouble)))
val indexed = points.zipWithIndex()
val indexedData = indexed.map{case (value, index) => (index, value)}
val pairedSamples = indexedData.cartesian(indexedData)
val dist = pairedSamples.map{case (x,y) => ((x,y),distance(x._2,y._2))}.map{case ((x,y),z) => (((x,y),z,covariance(z)))}    
val entries: RDD[MatrixEntry] = dist.map{case (((x,y),z,cov)) => MatrixEntry(x._1, y._1, cov)}
val coomat: CoordinateMatrix = new CoordinateMatrix(entries)        

进一步说明,我在 Spark 中从 RDD 创建了这个矩阵。那么也许将数据从 RDD 保存到 Mongodb 会更好/可能?

【问题讨论】:

    标签: mongodb matrix apache-spark rdd


    【解决方案1】:

    CoordinateMatrixRowMatrix 基本上分别是RDD[MatrixEntry]RDD[Vector] 的包装器,两者都可以相对保存到MongoDB。对于坐标矩阵:

    val spark: SparkSession = ???
    import spark.implicits._
    
    // For 1.x
    // val sqlContext: SQLContext = ???
    // import sqlContext.implicits._
    
    val options = Map(
       "uri" -> ???
       "database" -> ???
    )
    
    val coordMat = new CoordinateMatrix(sc.parallelize(Seq(
      MatrixEntry(1, 3, 1.4), MatrixEntry(3, 6, 2.8))
    ))
    
    coordMat.entries.toDF().write
      .options(options)
      .option("collection", "coordinates")    
      .format("com.mongodb.spark.sql")
      .save()
    

    你会得到形状的文件:

    {'_id': ObjectId('...'), 'i': 3, 'j': 6, 'value': 2.8}
    

    可以很容易地转换回原来的形式:

    val entries = spark.read
      .options(options)
      .option("collection", "coordinates")    
      .format("com.mongodb.spark.sql")
      .load()
      .drop("_id")  
      .schema(...)
      .as[MatrixEntry]
    
    new CoordinateMatrix(entries.rdd)
    

    几乎可以对RowMatrix 做同样的事情,但您需要做更多的工作(将Vectors 表示为密集数组或稀疏元组(size, indices, values))。

    不幸的是,在这两种情况下(CoordinateMatrixRowMatrix),您都会丢失有关矩阵形状的信息。

    【讨论】:

    • 谢谢。我收到此错误:“值 toDF 不是 org.apache.spark.rdd.RDD[....MatrixEntry] 的成员
    • import spark.implicits._ 其中sparkSparkSession 对象。对于 1.x,请使用 SQLContext
    • 是的。我添加了“import org.apache.spark.sql.SQLContext._”并将 coorMat 结构移到主类之外,但问题仍然存在。
    • 导入不正确。您需要来自活动实例的隐式。
    • Ahh .. 切换到 mongo-spark-connector 的版本 1.0.0,它现在可以工作了。我正在使用2.0.0-rc0。我还需要将spark.mongodb.output.collection 添加到 SparkConf。我认为添加到您的帖子中会很好。感谢您的帮助。
    猜你喜欢
    • 1970-01-01
    • 2015-04-23
    • 2014-07-31
    • 1970-01-01
    • 1970-01-01
    • 2011-08-12
    • 1970-01-01
    • 1970-01-01
    • 2014-05-17
    相关资源
    最近更新 更多