Spark是否自动复用相同DataFrame的Join执行计划以避免重复计算?
Spark Join结果复用问题解答
问题场景代码
def join_one(df_1, df_2): df = df_1.alias("d1").join(df_2.alias("d2")).select("d1.InvoiceNo","d1.StockCode","d1.Description","d2.CustomerID","d2.Country") return df def join_two(df_1, df_2): df = df_1.alias("d1").join(df_2.alias("d2")).select("d1.InvoiceDate","d1.UnitPrice","d2.Quantity") return df df_1 = spark.read\ .option("header", "true")\ .csv(path_1) df_2 = spark.read\ .option("header", "true")\ .csv(path_2) df_3 = join_one(df_1, df_2) df_4 = join_two(df_1, df_2) df_3.show() df_4.show()
问题描述
上述代码中,两次使用相同的df_1和df_2执行Join操作后选取不同列,想问Spark是否会自动优化执行计划,复用同一次Join的结果再选取对应列?还是会执行两次Join操作,需要用户自行优化?
解答
Spark不会自动复用这次Join的结果,会执行两次独立的Join操作。
原因是:Spark的Catalyst优化器虽然支持多种逻辑优化,但在这个场景下,df_3和df_4属于两个独立的DataFrame血统,它们的转换逻辑是分开定义的,优化器不会主动将这两个操作合并为一次Join后再拆分选列。
如果要避免重复Join、节省计算资源,你可以手动优化:先执行一次完整的Join得到结果DataFrame,再基于这个结果分别选取需要的列,示例代码如下:
# 先执行一次完整的Join操作 joined_df = df_1.alias("d1").join(df_2.alias("d2")) # 从Join结果中分别选取所需列生成目标DataFrame df_3 = joined_df.select("d1.InvoiceNo","d1.StockCode","d1.Description","d2.CustomerID","d2.Country") df_4 = joined_df.select("d1.InvoiceDate","d1.UnitPrice","d2.Quantity") df_3.show() df_4.show()
这样Spark只会执行一次Join操作,后续仅做两次列选择,在数据量较大时能显著提升执行效率。
内容的提问来源于stack exchange,提问作者Tarique
相关产品推荐
相关产品推荐

