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
相关产品推荐
相关产品推荐

