【问题标题】:Spark performance: local faster than clusterSpark性能:本地比集群快
【发布时间】:2020-09-22 14:08:59
【问题描述】:

我尝试在我的家庭网络上设置一个 Spark 集群,但与独立运行相比,我没有看到任何性能提升 - 事实上,与我在本地运行时相比,它稍微慢了一点 [*]。有人可以帮忙/解释一下原因吗?

我做了如下:

  • 我正在使用 MovieLens 10M 数据集。 http://files.grouplens.org/datasets/movielens/ml-10m.zip
  • 对于我的本地集群,我有两台现代高性能 mac,它们的规格相同,每台 32 GB RAM
  • 当我只运行独立 Spark 实例时,即本地 [*] 我的时间是 30.82 分钟
  • 当我针对将两台 mac 连接到同一个 spark 集群的 spark 集群运行时,我的时间是 35 分钟

我的spark集群的spark-submit参数如下

spark-submit --class com.sundogsoftware.spark.MovieSimilarities10MDataset --deploy-mode cluster --master spark://mbp2.lan:7077  --driver-memory 1g
--num-executors 2 --executor-cores 8 --executor-memory 28g

这导致在每台 Mac 上使用 8 个内核并使用 28GB RAM(通过 spark webapp 确认)

当然,我希望给定两台硬件规格相同的机器并将它们连接到同一个 spark 集群,我会看到使用相同数据集 (MovieLens 10M) 的性能提高

不胜感激任何建议。我已经多次调整了执行器/内核/内存的数量,但没有任何影响。

谢谢

类如下:

package com.sundogsoftware.spark

import java.util.concurrent.TimeUnit

import org.apache.log4j._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.{IntegerType, LongType, StringType, StructType}
import org.apache.spark.sql.{Dataset, SparkSession}

// To run on EMR successfully + output results for Star Wars:
// aws s3 cp s3://sundog-spark/MovieSimilarities1MDataset.jar ./
// aws s3 cp s3://sundog-spark/ml-10M100K/movies.dat ./
// spark-submit --executor-memory 1g MovieSimilarities1MDataset.jar 260

object MovieSimilarities10MDataset {

  case class Movies(userID: Int, movieID: Int, rating: Int, timestamp: Long)
  case class MoviesNames(movieID: Int, movieTitle: String)
  case class MoviePairs(movie1: Int, movie2: Int, rating1: Int, rating2: Int)
  case class MoviePairsSimilarity(movie1: Int, movie2: Int, score: Double, numPairs: Long)

  def computeCosineSimilarity(spark: SparkSession, data: Dataset[MoviePairs]): Dataset[MoviePairsSimilarity] = {
    // Compute xx, xy and yy columns
    val pairScores = data
      .withColumn("xx", col("rating1") * col("rating1"))
      .withColumn("yy", col("rating2") * col("rating2"))
      .withColumn("xy", col("rating1") * col("rating2"))

    // Compute numerator, denominator and numPairs columns
    val calculateSimilarity = pairScores
      .groupBy("movie1", "movie2")
      .agg(
        sum(col("xy")).alias("numerator"),
        (sqrt(sum(col("xx"))) * sqrt(sum(col("yy")))).alias("denominator"),
        count(col("xy")).alias("numPairs")
      )

    // Calculate score and select only needed columns (movie1, movie2, score, numPairs)
    import spark.implicits._
    val result = calculateSimilarity
      .withColumn("score",
        when(col("denominator") =!= 0, col("numerator") / col("denominator"))
          .otherwise(null)
      ).select("movie1", "movie2", "score", "numPairs").as[MoviePairsSimilarity]

    result
  }

  /** Get movie name by given movie id */
  def getMovieName(movieNames: Dataset[MoviesNames], movieId: Int): String = {
    val result = movieNames.filter(col("movieID") === movieId)
      .select("movieTitle").collect()(0)

    result(0).toString
  }
  /** Our main function where the action happens */
  def main(args: Array[String]) {

    // Set the log level to only print errors
    Logger.getLogger("org").setLevel(Level.ERROR)

    val startTime = System.nanoTime

    // Create a SparkSession without specifying master
    val spark = SparkSession
      .builder
      .appName("MovieSimilarities10M")
      .getOrCreate()

    // Create schema when reading u.item
    val moviesNamesSchema = new StructType()
      .add("movieID", IntegerType, nullable = true)
      .add("movieTitle", StringType, nullable = true)

    // Create schema when reading u.data
    val moviesSchema = new StructType()
      .add("userID", IntegerType, nullable = true)
      .add("movieID", IntegerType, nullable = true)
      .add("rating", IntegerType, nullable = true)
      .add("timestamp", LongType, nullable = true)

    println("\nLoading movie names...")
    import spark.implicits._
    // Create a broadcast dataset of movieID and movieTitle.
    // Apply ISO-885901 charset
    val movieNames = spark.read
      .option("sep", "::")
      .option("charset", "ISO-8859-1")
      .schema(moviesNamesSchema)
      .csv("movies.dat")
      .as[MoviesNames]

    // Load up movie data as dataset
    val movies = spark.read
      .option("sep", "::")
      .schema(moviesSchema)
      .csv("ratings.dat")
      .as[Movies]

    val ratings = movies.select("userId", "movieId", "rating")

    // Emit every movie rated together by the same user.
    // Self-join to find every combination.
    // Select movie pairs and rating pairs
    val moviePairs = ratings.as("ratings1")
      .join(ratings.as("ratings2"), $"ratings1.userId" === $"ratings2.userId" && $"ratings1.movieId" < $"ratings2.movieId")
      .select($"ratings1.movieId".alias("movie1"),
        $"ratings2.movieId".alias("movie2"),
        $"ratings1.rating".alias("rating1"),
        $"ratings2.rating".alias("rating2")
      ).repartition(100).as[MoviePairs]

    val moviePairSimilarities = computeCosineSimilarity(spark, moviePairs).cache()

    if (args.length > 0) {
      val scoreThreshold = 0.88
      val coOccurenceThreshold = 1000.0

      val movieID: Int = args(0).toInt

      // Filter for movies with this sim that are "good" as defined by
      // our quality thresholds above
      val filteredResults = moviePairSimilarities.filter(
        (col("movie1") === movieID || col("movie2") === movieID) &&
          col("score") > scoreThreshold && col("numPairs") > coOccurenceThreshold)

      // Sort by quality score.
      val results = filteredResults.sort(col("score").desc).take(50)

      println("\nTop 50 similar movies for " + getMovieName(movieNames, movieID))
      for (result <- results) {
        // Display the similarity result that isn't the movie we're looking at
        var similarMovieID = result.movie1
        if (similarMovieID == movieID) {
          similarMovieID = result.movie2
        }
        println(getMovieName(movieNames, similarMovieID) + "\tscore: " + result.score + "\tstrength: " + result.numPairs)
      }

      val stopTime = System.nanoTime
      val elapsedTime = (stopTime - startTime)


      // TimeUnit
      val convert = TimeUnit.SECONDS.convert(elapsedTime, TimeUnit.NANOSECONDS)

      //    System.out.println(convert + " seconds")
      println(s"elapsedTime sec=$convert")
    }
  }
}

【问题讨论】:

  • 您可以尝试使用更多的执行器,但执行器的核心/内存更少吗?比如:--num-executors 5,--executor-cores 3,--executor-memory 10g
  • 谢谢,现在运行这个 --driver-memory 2g --num-executors 5 --executor-cores 3 --executor-memory 10g
  • 很遗憾集群的完成时间没有变化;还是35分钟。谢谢
  • 只有在计算阶段比加载、反序列化和洗牌两个输入文件(我假设每个文件都是一个分区)更昂贵时,集群才会更快?在您的情况下,这些将是乘法和新列。由于加载和收集都是一次使用几个执行程序或驱动程序的限制。测量每个大逻辑块的时间将帮助您了解瓶颈在哪里

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


【解决方案1】:

您从互联网上下载了两次相同的数据集,而不是一次。这比下载一次需要更长的时间。

附带说明,您在此处进行数值处理。使用 GPU 并在 gpu 上进行体矩阵运算,您可能会获得更多的加速。

【讨论】:

  • 感谢文件是本地文件,没有被下载。我已经删除了评论,这样就不会混淆 re:S3 参考。这是示例中的代码注释(这是课程代码)不是我的
  • 您已经在每个节点上预先配置了它?几乎可以肯定,您的加入正在引起整个网络的混乱。
  • 是的,在每个节点上预先配置。虽然是局域网。关于我如何看到这一点的任何提示?谢谢
  • 看看火花控制台。您将看到执行计划、洗牌所花费的时间、洗牌所花费的数据等。此外,$"ratings1.movieId" &lt; $"ratings2.movieId" 比等值连接更昂贵。我不清楚你为什么需要它。
  • 分布式 Spark 会产生很多独立运行无法获得的开销:通信、通过线路编组数据的需要等。
【解决方案2】:
  • 分布式 spark 有很多开销,您无法在独立环境中获得:通信、需要通过线路编组数据等。我认为您需要两个以上的工作人员才能离开独立环境是合理的。
  • 确保数据的格式可以轻松地在工作人员之间拆分。考虑在工作流程的第一步加载 .dat 文件并将其保存为 Parquet 或至少 avro 格式,可能使用 snappy(但不是 gzip)压缩。

您将需要一个共享文件系统才能正常工作;对于两台机器 NFS 是您正在使用的机器,对吗?将(转换后的)源文件粘贴在本地 FS 上完全相同的路径中,并将共享 NFS 用于所有其他工作。这样:没有网络 IO 来加载数据。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-08-27
    • 1970-01-01
    • 1970-01-01
    • 2022-06-15
    • 2018-12-10
    • 1970-01-01
    • 1970-01-01
    • 2020-08-10
    相关资源
    最近更新 更多