Spark中Subquery与DataFrame过滤函数的性能对比及原因
问题
我运行了以下带子查询的Spark SQL语句:
val df = spark.sql("""select * from employeesTableTempview where dep_id in (select dep_id from departmentTableTempview)""") df.count()
同时我也用DataFrame函数式方式实现了相同逻辑:假设已将员工表和部门表读取为DataFrame,分别命名为empDF和DepDF,代码如下:
val depidList = DepDF.map(x=>x(0).string).collect().toList() val empdf2 = empDF.filter(col("dep_id").isin(depidList:_*)) empdf2.count
请问在这两种场景下,哪种方式性能更优?原因是什么?请基于Spark Scala帮我理解这两种场景。
性能对比与原因分析
- Spark SQL子查询方式性能更优,核心差异在于两者的执行逻辑和Spark优化机制的利用程度:
数据处理的分布式特性差异
- 函数式方式中,
collect()会把DepDF的所有数据拉取到Driver端本地内存,再将本地列表作为过滤条件下发到Executor。如果DepDF数据量较大,不仅会占用Driver大量内存,甚至可能触发OOM;同时大列表的网络传输也会增加额外开销。 - Spark SQL的子查询会被Catalyst优化器自动解析为分布式Join操作:若
DepDF数据量小,会自动执行Broadcast Hash Join(将小表广播到所有Executor,避免Shuffle);若数据量大则选择Shuffle Hash Join等合适的分布式策略,所有计算都在Executor集群中完成,不会把大量数据拉到Driver,规避了单点压力。
- 函数式方式中,
优化器的作用差异
- Spark SQL会经过Catalyst优化器的多轮优化:包括谓词下推、子查询转Join、广播自动选择等,能根据数据规模自动匹配最优执行计划。
- 函数式方式的
isin操作是基于本地列表的硬编码过滤,Spark无法对该逻辑进行分布式优化,只能严格按照用户编写的步骤执行,优化空间几乎为零。
数据量适应性差异
- 函数式方式仅适合
DepDF数据量极小的场景(如几百条以内),数据量增大后会迅速出现性能瓶颈甚至崩溃。 - Spark SQL子查询方式适配所有数据量级,无论
DepDF是小表还是大表,优化器都会选择对应的分布式执行方案,性能表现更稳定。
- 函数式方式仅适合
内容的提问来源于stack exchange,提问作者nirmal
相关产品推荐
相关产品推荐

