Spark如何获取广播判定用DataFrame大小?两个统计值哪个正确?
问题解答
1. Catalyst统计值是否为自动广播阈值的判断依据
你通过optimizedPlan.stats.sizeInBytes拿到的58MB就是Spark判断是否触发自动广播的参考值。Spark SQL的自动广播逻辑完全依赖Catalyst优化器输出的统计值做阈值比较,如果你要调整spark.sql.autoBroadcastJoinThreshold只针对bdrDf生效,直接把阈值设置为略高于该值即可。
2. 为什么Spark UI显示的大小和统计值不一致
两个值对应不同的计算场景,不存在谁对谁错:
- 代码拿到的58MB是逻辑统计大小:是Catalyst根据表元数据、字段类型、采样统计等估算的纯数据逻辑大小,只用于优化策略判断,不代表实际运行时的内存占用。
- Spark UI看到的492MB是实际缓存物理大小:是DataFrame缓存到内存后实际占用的内存空间,包含了Java对象头开销、序列化开销、存储结构开销。如果是未序列化缓存,Java对象的内存占用会远高于逻辑数据大小,属于正常现象。
另外你写的测试代码中,先缓存再拿逻辑统计值的操作是多余的,逻辑统计值只和DataFrame的查询计划有关,和是否缓存、是否触发执行没有关系。
3. 不同场景下的DataFrame大小估算方法
场景1:调整自动广播阈值、做SQL执行计划优化
直接使用你当前的统计逻辑即可,这也是Spark官方做优化决策的标准参考:
import org.apache.commons.io.FileUtils val bytes = sparkSession.sessionState.executePlan(bdrDf.queryExecution.logical) .optimizedPlan.stats(sparkSession.sessionState.conf).sizeInBytes println("逻辑大小:" + FileUtils.byteCountToDisplaySize(bytes.toLong))
场景2:估算实际内存占用、做资源规划
可以缓存DataFrame触发执行后,直接看Spark UI Storage标签的对应大小,也可以通过代码直接读取缓存的物理占用:
df.cache() df.count() // 触发全量缓存 // 读取实际缓存大小 spark.sharedState.cacheManager.lookupCachedData(df).foreach { cachedData => val physicalBytes = cachedData.cachedRepresentation.cacheBuilder.sizeInBytes println("实际内存占用:" + FileUtils.byteCountToDisplaySize(physicalBytes)) }
内容的提问来源于stack exchange,提问作者Marwan02
相关产品推荐
相关产品推荐

