Spark小DataFrame分区过多及广播Join相关技术问题咨询
问题解答
1. 少量数据多分区的原因、性能影响及创建方式正确性
原因
Spark从本地集合创建DataFrame时,分区数由spark.default.parallelism参数决定,Databricks集群会根据Executor的CPU核数自动计算该参数的值(通常为总核数的1-3倍),因此即使只有27行数据,也会按照集群默认并行度设置分区,出现数十个分区的情况。
性能影响
会带来明显性能损耗:每个分区对应一个任务,任务的启动、调度、序列化/反序列化都有固定开销,当数据量极小时,任务实际计算时间远小于调度开销,整体执行效率会下降。
创建方式正确性
原代码存在两处语法错误,修正后即可正常使用:
- Schema定义中,每个
StructField末尾缺少逗号,正确写法:
schema=StructType( [ StructField("col1",IntegerType(),True), StructField("col2",StringType(),True), StructField("col3",StringType(),True), StructField("col4",IntegerType(),True), StructField("col5",IntegerType(),True) ] )
- 数据列表中存在多余的空逗号(
,,),会触发语法报错,需删除这些无效逗号,仅保留有效数据行的分隔逗号。
修正后的创建方式语法合规,可正常生成DataFrame。
2. 未触发BroadcastJoin的原因
主要有以下几种可能:
- 缺失统计信息:从本地集合创建的DataFrame默认未收集表统计数据,Spark优化器无法准确判断其实际大小,可能误判为超过广播阈值,从而选择SortMergeJoin。可通过
ANALYZE TABLE mydf COMPUTE STATISTICS命令手动收集统计信息,帮助优化器做出正确判断。 - 广播阈值被隐性修改:虽然你未手动设置
spark.sql.autoBroadcastJoinThreshold = -1,但Databricks集群的默认配置、初始化脚本或会话级配置可能已将该阈值调至过小,导致27行数据的序列化大小超过阈值(比如字符串字段实际内容较长时)。 - 关联键特殊性:如果关联键包含大量
null值,Spark优化器可能认为BroadcastJoin执行效率不如SortMergeJoin,但这种情况较少见。
3. 查看BroadcastJoin默认阈值的命令
在Databricks中可通过两种方式查看:
- PySpark命令:
spark.conf.get("spark.sql.autoBroadcastJoinThreshold")
- SQL命令:
SET spark.sql.autoBroadcastJoinThreshold;
内容的提问来源于stack exchange,提问作者Cassius Clay
相关产品推荐
相关产品推荐

