Spark中基于JOIN缩小表范围的实现及性能优化咨询
问题解答
方案可行性
你提供的临时表+SQL方案完全可以实现需求:通过EXISTS子查询能精准筛选出ratings中存在的影片记录,同时直接全量保存ratings数据即可满足要求。
性能优化措施
针对1000万条movies和100万条ratings的场景,可从以下几个维度优化性能:
1. 利用广播小表消除Shuffle
ratings属于小表(仅100万条),通过广播小表可避免大表(movies)与小表的全量Shuffle操作,这是最有效的优化手段:
- 方式一:在SQL中添加广播提示
SELECT titles.movie_id, titles.movie_type, titles.primary_title, titles.original_title, titles.is_adult, titles.start_year, titles.end_year, titles.runtime_minutes, titles.genres FROM TMPtitles titles WHERE EXISTS (SELECT 1 FROM /*+ BROADCAST(TMPratings) */ TMPratings ratings WHERE titles.movie_id = ratings.movie_id) - 方式二:先提取
ratings的唯一movie_id并广播,再关联val uniqueMovieIds = rating.select("movie_id").distinct() import org.apache.spark.sql.functions.broadcast val filteredMovies = title.join(broadcast(uniqueMovieIds), Seq("movie_id"), "inner")
2. 提前精简小表数据
如果ratings中存在重复的movie_id,先对其去重后再创建临时表,减少后续广播和关联的数据量:
val distinctRatings = rating.select("movie_id").distinct() distinctRatings.createOrReplaceTempView("TMPratings")
3. 指定Schema避免自动推断
读取S3上的原始数据时,手动指定Schema,避免Spark为了推断Schema而进行全表扫描,节省时间:
import org.apache.spark.sql.types._ val movieSchema = StructType(Seq( StructField("movie_id", StringType), StructField("movie_type", StringType), // 其他字段依次定义对应类型 )) val title = spark.read.schema(movieSchema).csv("s3://your-path/movies")
4. 调整Spark执行配置
- 开启自适应执行:让Spark根据运行时数据自动调整执行计划
spark.sql.adaptive.enabled true - 调整Shuffle分区数:默认200,可根据集群资源调整(比如设置为Executor核数的2倍)
spark.sql.shuffle.partitions 100 - 合理分配Executor内存与核数:根据集群资源提升并行处理能力
5. 替换EXISTS为内连接(可选)
Spark优化器对INNER JOIN的优化支持成熟,结合广播提示,性能与EXISTS相当甚至更优,逻辑也更直观:
SELECT DISTINCT titles.movie_id, titles.movie_type, titles.primary_title, titles.original_title, titles.is_adult, titles.start_year, titles.end_year, titles.runtime_minutes, titles.genres FROM TMPtitles titles INNER JOIN /*+ BROADCAST(TMPratings) */ TMPratings ratings ON titles.movie_id = ratings.movie_id
6. S3读取优化
- 使用S3A文件系统替代旧的S3实现,提升S3读写性能:
spark.hadoop.fs.s3.impl org.apache.hadoop.fs.s3a.S3AFileSystem - 调整S3块大小与批量读取参数,减少网络IO开销:
spark.hadoop.fs.s3a.block.size 134217728 -- 设置为128MB spark.hadoop.fs.s3a.fast.upload true
内容的提问来源于stack exchange,提问作者Hana
相关产品推荐
相关产品推荐

