You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 02:36:03