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

如何用PySpark基于节点与关联DataFrame动态生成链路关系表

PySpark实现方案

核心实现思路:

  • 先通过三次关联把DF2每行的三个角色名称都映射为DF1中对应的ID,得到包含所有ID的中间宽表
  • 再将每行对应的两段链路拆分为独立行,最终得到目标边表

具体可运行代码如下:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

# 初始化SparkSession
spark = SparkSession.builder.appName("generate_edge_table").getOrCreate()

# ---------------------- 示例数据构造(实际场景替换为你自己的读表逻辑即可) ----------------------
# 节点表DF1
df1_data = [("A",0,"mgr"),("B",1,"mgr"),("C",2,"mgr"),
            ("D",3,"hr"),("E",4,"hr"),("F",5,"hr"),
            ("G",6,"adm"),("H",7,"adm"),("I",8,"adm")]
df1 = spark.createDataFrame(df1_data, schema=["Name", "ID", "Group"])

# 关联关系表DF2
df2_data = [("A","D","G",0.0010),("B","E","H",0.0002),("C","F","I",0.0035)]
df2 = spark.createDataFrame(df2_data, schema=["Mgrs", "HR", "Admin", "Value"])

# ---------------------- 核心处理逻辑 ----------------------
# 三次关联映射三个角色对应的ID
mid_df = df2.join(df1.select("Name", "ID").alias("mgr"), col("Mgrs") == col("mgr.Name"), "left") \
            .withColumnRenamed("ID", "mgr_id") \
            .join(df1.select("Name", "ID").alias("hr"), col("HR") == col("hr.Name"), "left") \
            .withColumnRenamed("ID", "hr_id") \
            .join(df1.select("Name", "ID").alias("adm"), col("Admin") == col("adm.Name"), "left") \
            .withColumnRenamed("ID", "adm_id") \
            .select("mgr_id", "hr_id", "adm_id", "Value")

# 拆分行生成两段链路,两种方案二选一即可
# 方案1:union拼接,逻辑清晰易读
df3 = mid_df.select(col("mgr_id").alias("From"), col("hr_id").alias("To"), col("Value")) \
            .unionAll(
                mid_df.select(col("hr_id").alias("From"), col("adm_id").alias("To"), col("Value"))
            )

# 方案2:stack函数拆分,大数据量下性能更优
# df3 = mid_df.selectExpr("stack(2, mgr_id, hr_id, hr_id, adm_id) as (From, To)", "Value")

# 如需和示例输出顺序一致,可追加排序操作
# df3 = df3.orderBy("From")

df3.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 10:57:00