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

PySpark DataFrame为每个ID添加最新ID列的实现方法求助

解决PySpark中ID链的最新ID映射问题

这个问题本质是处理ID的连通链:每个旧ID更新为新ID,形成一条链式关系,自连接只能处理单次变更,无法覆盖多级链,因此需要用图连通分量或递归CTE的方式解决。

方法一:使用GraphFrames处理连通分量

GraphFrames可以轻松识别出所有属于同一ID链的节点,再找到每个链中最新的ID。

步骤1:初始化环境与示例数据

from pyspark.sql import SparkSession
from pyspark.sql.functions import max, col
from graphframes import GraphFrame

spark = SparkSession.builder.appName("LatestIDMapping").getOrCreate()

# 构建示例DataFrame
data = [
    ("d", "c", "20-01-01"),
    ("c", "b", "15-01-01"),
    ("b", "a", "10-01-01"),
    ("z", "y", "23-02-01"),
    ("y", "x", "20-01-02")
]
df = spark.createDataFrame(data, ["new_id", "old_id", "change_date"])

步骤2:构建图结构

  • 顶点:所有出现过的ID(包含new_id和old_id的唯一值)
  • 边:表示ID的更新关系(从old_id指向new_id)
# 生成顶点表
vertices = df.selectExpr("new_id as id").union(df.selectExpr("old_id as id")).distinct()

# 生成边表
edges = df.selectExpr("old_id as src", "new_id as dst")

步骤3:计算连通分量

同一个ID链的所有节点会被分配同一个component ID:

# 创建GraphFrame并计算连通分量
g = GraphFrame(vertices, edges)
connected_components = g.connectedComponents()

步骤4:找到每个分量的最新ID

通过分组找到每个连通分量中change_date最大的记录,对应的new_id就是该链的最新ID:

# 将原数据与连通分量关联
df_with_component = df.join(connected_components, df.new_id == connected_components.id, "left")

# 提取每个分量的最新ID
latest_per_component = df_with_component.groupBy("component")\
    .agg(max("change_date").alias("max_date"))\
    .join(df_with_component, (col("component") == df_with_component.component) & (col("max_date") == df_with_component.change_date))\
    .select("component", "new_id as latest")

步骤5:关联回原DataFrame

将最新ID映射到原数据的每条记录:

# 关联分量ID与最新ID,得到最终结果
result = df.join(connected_components, df.old_id == connected_components.id, "left")\
    .join(latest_per_component, "component", "left")\
    .select("new_id", "old_id", "change_date", "latest")

result.show()

运行后会输出你期望的结果。

方法二:使用递归CTE(Spark 3.0+支持)

如果你的Spark版本≥3.0,可以用SQL递归CTE遍历ID链,直接找到每个记录的最新ID:

# 注册临时表
df.createOrReplaceTempView("id_changes")

# 执行递归CTE
result = spark.sql("""
WITH RECURSIVE id_chain AS (
    -- 基础节点:没有后续更新的ID(new_id不在old_id列表中)
    SELECT new_id, old_id, change_date, new_id AS latest
    FROM id_changes
    WHERE new_id NOT IN (SELECT old_id FROM id_changes)
    UNION ALL
    -- 递归遍历:找到指向当前链的旧ID记录
    SELECT d.new_id, d.old_id, d.change_date, ic.latest
    FROM id_changes d
    JOIN id_chain ic ON d.new_id = ic.old_id
)
SELECT * FROM id_chain ORDER BY change_date DESC
""")

result.show()

这种方法无需额外依赖GraphFrames,适合纯SQL场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 03:25:00