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

PySpark(Spark2.2)左连接聚合MAX取分组最新记录实现方法

PySpark 2.2 按IG关联取最新服务记录实现方案

适用于部署在Hadoop上的Spark 2.2 PySpark环境,针对大规模DataFrame做了性能优化,避免Join阶段数据膨胀。

核心思路

直接全表Join后再分组取最大值会导致中间数据量暴涨,大表场景下极易出现Executor OOM,因此优先对DataFrame B做预处理:按IG分组后筛选出每个IG下Dt_Service最大(最新服务时间)的唯一记录,再和DataFrame A做Left Join,性能更稳定。

实现代码

首先导入依赖:

from pyspark.sql import Window
import pyspark.sql.functions as F

可以用下面的测试样例先验证逻辑,后续替换成自己的真实表即可:

# 测试用DataFrame A 样例数据
data_a = [
    (1, "IG001", "2023-01-05"),
    (2, "IG001", "2023-02-10"),
    (3, "IG002", "2023-03-15"),
    (4, "IG003", "2023-04-20")
]
df_a = spark.createDataFrame(data_a, schema=["ID", "IG", "OpenDate"])

# 测试用DataFrame B 样例数据
data_b = [
    ("IG001", "ServiceA", "2023-01-01"),
    ("IG001", "ServiceB", "2023-02-15"),
    ("IG001", "ServiceC", "2023-01-20"),
    ("IG002", "ServiceD", "2023-03-10"),
    ("IG002", "ServiceE", "2023-03-20")
]
df_b = spark.createDataFrame(data_b, schema=["IG", "Service", "Dt_Service"])

核心处理逻辑:

# 定义窗口分区规则:按IG分组,组内按服务时间倒序
window_rule = Window.partitionBy("IG").orderBy(F.desc("Dt_Service"))

# 预处理B表,仅保留每个IG下最新的服务记录
df_b_latest = df_b.withColumn("row_num", F.row_number().over(window_rule)) \
    .filter(F.col("row_num") == 1) \
    .select("IG", "Service", "Dt_Service")

# 左关联得到最终结果
df_result = df_a.join(df_b_latest, on="IG", how="left") \
    .select("ID", "IG", "OpenDate", "Service", "Dt_Service")

执行查看结果:

df_result.show()

# 输出结果完全匹配期望字段:ID、IG、OpenDate、Service、Dt_Service
#+---+-----+----------+--------+----------+
#| ID|   IG|  OpenDate|Service|Dt_Service|
#+---+-----+----------+--------+----------+
#|  1|IG001|2023-01-05|ServiceB|2023-02-15|
#|  2|IG001|2023-02-10|ServiceB|2023-02-15|
#|  3|IG002|2023-03-15|ServiceE|2023-03-20|
#|  4|IG003|2023-04-20|    null|      null|
#+---+-----+----------+--------+----------+

注意事项

  • 如果同一个IG下存在多条记录的Dt_Service同为最大值,row_number()会随机返回其中一条;如果需要保留所有同最大时间的记录,把row_number()替换为rank()即可。
  • 如果Dt_Service字段是字符串格式,必须保证是yyyy-MM-dd、yyyy-MM-dd HH:mm:ss这类按字典序排序和时间顺序一致的格式,否则需要先通过F.to_date()/F.to_timestamp()转成日期时间类型再排序,避免取到错误的最新记录。
  • 预处理B表的逻辑会做Map端预聚合,Shuffle数据量远小于先Join再分组的方案,适合亿级以上的大规模表关联场景。

内容的提问来源于stack exchange,提问作者stay_funn

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 22:15:30