【发布时间】: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