【问题标题】:How to distribute code and dataset onto worker nodes?如何将代码和数据集分发到工作节点?
【发布时间】:2019-11-11 01:01:01
【问题描述】:

我一直在使用数据集 Movielens(2000 万条记录),并且一直在 Spark MLlib 中使用 collaborative filtering

我的环境是 VirtualBox 上的 Ubuntu 14.4。我有一个主节点和两个从节点。我使用了已发布的 Apache Hadoop、Apache Spark、Scala、sbt。代码是用 Scala 编写的。

如何将代码和数据集分发到工作节点上?

import java.lang.Math._

import org.apache.spark.ml.recommendation.ALS
import org.apache.spark.ml.recommendation.ALS.Rating
import org.apache.spark.sql.SQLContext
import org.apache.spark.{SparkConf, SparkContext}

object trainModel extends App {

  val conf = new SparkConf()
    .setMaster("local[*]")
    .setAppName("trainModel")
  val sc = new SparkContext(conf)

  val rawData = sc.textFile("file:///usr/local/spark/dataset/rating.csv")

  val sqlContext = new SQLContext(sc)
  val df = sqlContext
    .read
    .option("header", "true")
    .format("csv")
    .load("file:///usr/local/spark/dataset/rating.csv")

  val ratings = rawData.map(line => line.split(",").take(3) match {
    case Array(userId, movieId, rating) => 
      Rating(userId.toInt, movieId.toInt, rating.toFloat)
  })
  println(s"Number of Ratings in Movie file ${ratings.count()} \n")

  val ratingsRDD = sc.textFile("file:///usr/local/spark/dataset/rating.csv")
  //split data into test&train
  val splits = ratingsRDD.randomSplit(Array(0.8, 0.2), seed = 12345)
  val trainingRatingsRDD = splits(0).cache()
  val testRatingsRDD = splits(1).cache()
  val numTraining = trainingRatingsRDD.count()
  val numTest = testRatingsRDD.count()
  println(s"Training: $numTraining, test: $numTest.")

  val rank = 10
  val lambdas = 0.01
  val numIterations = 10
  val model = ALS.train(ratings, rank, numIterations)
  //Evaluate the model on training data
  val userProducts = ratings.map { case Rating(userId, movieId, rating) =>
    (userId, movieId)
  }
  val predictions = model.predict(userProducts).map { case
    Rating(userId, movieId, rating) =>
    ((userId, movieId), rating)
  }
  val ratesAndPreds = ratings.map { case Rating(userId, movieId, rating) =>
    ((userId, movieId),
      rating)
  }.join(predictions)
  val meanSquaredError = ratesAndPreds.map { case ((userId, movieId),
  (r1, r2)) =>
    val err = r1 - r2
    err * err
  }.mean
  println("Mean Squared Error= " + meanSquaredError)
  sqrt(meanSquaredError)
  val rmse = math.sqrt(meanSquaredError)
  println(s" RMSE = $rmse.")
}

【问题讨论】:

    标签: scala apache-spark apache-spark-sql apache-spark-mllib


    【解决方案1】:

    如何分发代码

    当您spark-submit 一个 Spark 应用程序时会发生这种情况。分配可以按 CPU 内核/线程或执行程序进行。您不必对其进行编码。这就是人们使用 Spark 的原因,因为它应该(几乎)自动发生。

    conf.setMaster("local[*]")

    也就是说,您使用的单个执行程序的线程数与您拥有的 CPU 内核数一样多。这是一个本地分布。

    您最好从代码中删除该行并改用spark-submit --master。阅读官方文档,尤其是。 Submitting Applications.

    ...和数据集到工作节点? val rawData = sc.textFile("file:///usr/local/spark/dataset/rating.csv")

    该行说明了 Movielens 数据集 (rating.csv) 的分布方式。它与 Spark 无关,因为 Spark 使用文件系统上的任何分布。

    换句话说,在具有 256MB 块大小 (split) 的 Hadoop HDFS 上,一个两倍于块大小的文件可分为两部分。那就是 HDFS 使文件分布式和容错。

    当 Spark 读取 2-split 文件时,分布式计算(使用 RDD 描述)将使用 2 个分区和 2 个任务。

    HDFS 是一个文件系统/存储,所以选择任何位置和hdfs -put 数据集。将 HDFS 视为您可以远程访问的任何文件系统。使用位置作为sc.textFile 的输入参数,就完成了。

    【讨论】:

    • core-site.xml 文件是 fs.defaultFShdfs://0.0.0.0:9000hadoop.tmp.dir/app/hadoop/tmp
    • 您是否使用hdfs -put 命令上传数据集?请这样做,并将文件上传到您选择的位置。将其用作sc.textFile 的输入参数即可。
    • 非常感谢您的解释和您的努力。
    • 如果对您有用,请接受答案。谢谢。
    • 请问,我可以用 eclipse 和 scala IDE 来运行它吗?
    【解决方案2】:

    1 - 您的数据集最好放置在分布式文件系统中 - Hadoop HDFS、S3 等。

    2 - 代码通过spark-submit 脚本分发,如此处所述https://spark.apache.org/docs/2.4.3/submitting-applications.html

    【讨论】:

    • 拜托,你能详细解释一下吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-01-29
    • 1970-01-01
    • 2021-01-12
    • 2020-04-23
    • 2016-06-15
    • 1970-01-01
    • 2017-07-15
    相关资源
    最近更新 更多