PySpark如何对单表执行GROUP BY后与另一张表关联查询
原代码错误原因
- 运算符优先级问题:PySpark中
&的优先级高于==,你写的判断条件没有加括号包裹等式,会触发语法错误 - 逻辑错误:join条件中直接调用
T2.groupBy('Common').agg(f.max('Value'))返回的是包含Common和max(Value)两列的聚合DataFrame,不能直接和T2.Value单个列做等值判断,不符合语法规则
实现方案1:窗口函数实现(推荐,性能更优)
这种方式只需要扫描一次Table2,即可拿到每个Common分组下Value最大的行,再和Table1关联即可
import pyspark.sql.functions as f from pyspark.sql.window import Window # 定义窗口按Common分组,Value倒序排序 w = Window.partitionBy("Common").orderBy(f.col("Value").desc()) # 给Table2加行号,取每个分组第一行(也就是Value最大的行) table2_max = Table2.withColumn("rn", f.row_number().over(w)) \ .filter(f.col("rn") == 1) \ .drop("rn") # 和Table1左关联得到最终结果 result = Table1.alias("t1").join( table2_max.alias("t2"), on="Common", how="left" )
如果需要保留同一Common下所有Value等于最大值的行,可以把row_number()替换为rank()
实现方案2:与你提供的子查询SQL逻辑完全对齐
import pyspark.sql.functions as f # 先得到每个Common的最大Value t2_max_val = Table2.groupBy("Common").agg(f.max("Value").alias("max_val")) # 关联原Table2过滤出每个分组Value最大的行 table2_max = Table2.alias("t2").join( t2_max_val.alias("max_t2"), (f.col("t2.Common") == f.col("max_t2.Common")) & (f.col("t2.Value") == f.col("max_t2.max_val")), how="inner" ).select("t2.*") # 和Table1左关联得到最终结果 result = Table1.alias("t1").join( table2_max.alias("t2"), on="Common", how="left" )
内容的提问来源于stack exchange,提问作者Debtanu Gupta
相关产品推荐
相关产品推荐

