Spark 3.2中获取DataFrame大小的原有方法失效求助
Spark 3.2+ 中获取DataFrame预估字节大小的解决方法
Spark 3.2 版本起,SessionState.executePlan() 方法签名发生变更,要求传入逻辑计划和命令执行模式两个参数,原仅传逻辑计划的调用方式会触发方法不存在的错误。以下是两种可行的解决方式:
方式一:适配新的executePlan签名
通过Py4J正确获取CommandExecutionMode枚举实例,再调用变更后的executePlan方法:
# 获取执行模式枚举实例 execution_mode = spark_session._jsparkSession.CommandExecutionMode.EXECUTE_VALUE() # 计算DataFrame预估字节大小 size_in_bytes = spark_session._jsparkSession.sessionState().executePlan( df._jdf.queryExecution().logical(), execution_mode ).optimizedPlan().stats().sizeInBytes()
关键说明:
- 不能直接访问
spark_session._jsparkSession.CommandExecutionMode属性,需通过EXECUTE_VALUE()方法获取枚举实例,否则会触发AttributeError。 EXECUTE_VALUE是Spark默认的执行模式,对应Scala端的CommandExecutionMode.ExecuteValue。
方式二:直接从DataFrame查询计划获取
无需通过SessionState,直接从DataFrame的查询执行计划中提取优化后的统计信息:
size_in_bytes = df._jdf.queryExecution().optimizedPlan().stats().sizeInBytes()
补充说明:
如果返回值为-1,说明统计信息未预计算,可手动触发统计更新:
# 创建临时视图并触发统计计算 df.createOrReplaceTempView("temp_df") spark_session.sql("ANALYZE TABLE temp_df COMPUTE STATISTICS") # 重新获取大小 size_in_bytes = df._jdf.queryExecution().optimizedPlan().stats().sizeInBytes()
内容的提问来源于stack exchange,提问作者nadavsh
相关产品推荐
相关产品推荐

