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

如何以Spark高效方式协调DataFrame列的行值?

Spark DataFrame跨EventId协调ID字段的高效实现方案

原始DataFrame定义

val df1 = Seq(
    ("EventId1", Some("GUID1"), Some("ID1"),    None),
    ("EventId2", None,          Some("ID1"),    Some("Uid1")),
    ("EventId3", Some("GUID1"), None,           None),
    ("EventId4", Some("GUID3"), Some("ID3"),    None),
    ("EventId5", None,          Some("ID3"),    Some("Uid3"))
).toDF("EventId", "GUID", "WID", "SUid")

原始数据展示:

+--------+-----+----+----+
| EventId| GUID| WID|SUid|
+--------+-----+----+----+
|EventId1|GUID1| ID1|null|
|EventId2| null| ID1|Uid1|
|EventId3|GUID1|null|null|
|EventId4|GUID3| ID3|null|
|EventId5| null| ID3|Uid3|
+--------+-----+----+----+

需求说明

跨不同EventId协调补全GUID、WID、SUid三个字段:属于同一关联组的记录(如GUID1与ID1关联、GUID3与ID3关联),这三个字段需统一为该组的完整非空值。

预期结果

+--------+-----+---+----+
| EventId| GUID|WID|SUid|
+--------+-----+---+----+
|EventId1|GUID1|ID1|Uid1|
|EventId2|GUID1|ID1|Uid1|
|EventId3|GUID1|ID1|Uid1|
|EventId4|GUID3|ID3|Uid3|
|EventId5|GUID3|ID3|Uid3|
+--------+-----+---+----+

实现方案

方案一:分组关联构建映射表(无需额外依赖)

核心思路是先提取各组的完整ID映射关系,再与原表关联补全字段,适配大多数常规场景:

import org.apache.spark.sql.functions._

// 1. 按WID分组,提取每个WID对应的非空GUID和SUid
val wid_mapping = df1.filter(col("WID").isNotNull)
  .groupBy("WID")
  .agg(
    first(when(col("GUID").isNotNull, col("GUID"))).alias("group_guid"),
    first(when(col("SUid").isNotNull, col("SUid"))).alias("group_suid")
  )

// 2. 按GUID分组,关联WID映射表补全WID和SUid信息
val guid_mapping = df1.filter(col("GUID").isNotNull)
  .join(wid_mapping, df1("WID") === wid_mapping("WID"), "left")
  .select(
    col("GUID").alias("group_guid"),
    coalesce(df1("WID"), wid_mapping("WID")).alias("group_wid"),
    wid_mapping("group_suid")
  )
  .groupBy("group_guid")
  .agg(
    first(when(col("group_wid").isNotNull, col("group_wid"))).alias("group_wid"),
    first(when(col("group_suid").isNotNull, col("group_suid"))).alias("group_suid")
  )

// 3. 合并两类映射表,得到完整的组ID映射
val full_mapping = wid_mapping.join(guid_mapping, wid_mapping("group_guid") === guid_mapping("group_guid"), "full")
  .select(
    coalesce(wid_mapping("WID"), guid_mapping("group_wid")).alias("WID"),
    coalesce(wid_mapping("group_guid"), guid_mapping("group_guid")).alias("GUID"),
    coalesce(wid_mapping("group_suid"), guid_mapping("group_suid")).alias("SUid")
  )
  .distinct()

// 4. 关联原表补全所有记录的ID字段
val result_df = df1.join(
  full_mapping,
  (df1("GUID") === full_mapping("GUID")) || (df1("WID") === full_mapping("WID")),
  "left"
).select(
  df1("EventId"),
  full_mapping("GUID"),
  full_mapping("WID"),
  full_mapping("SUid")
)

result_df.show()

方案二:图计算连通组件(适配复杂关联场景)

如果ID存在多层关联等复杂关系,可通过GraphFrames找出连通组件,同一组件内的记录共享ID信息:

import org.graphframes.GraphFrame
import org.apache.spark.sql.functions._

// 1. 创建顶点表:所有非空的GUID和WID作为节点
val vertices = df1.select(col("GUID").alias("id")).filter(col("id").isNotNull)
  .union(df1.select(col("WID").alias("id")).filter(col("id").isNotNull))
  .distinct()

// 2. 创建边表:建立GUID与WID之间的双向关联边
val edges = df1.filter(col("GUID").isNotNull && col("WID").isNotNull)
  .select(col("GUID").alias("src"), col("WID").alias("dst"))
  .union(df1.filter(col("GUID").isNotNull && col("WID").isNotNull)
    .select(col("WID").alias("src"), col("GUID").alias("dst")))

// 3. 构建图并计算连通组件
val graph = GraphFrame(vertices, edges)
val components = graph.connectedComponents.run()

// 4. 按连通组件分组,提取每组的完整ID信息
val component_info = df1.join(components, df1("GUID") === components("id"), "left")
  .join(components.withColumnRenamed("id", "wid_id").withColumnRenamed("component", "component2"), df1("WID") === col("wid_id"), "left")
  .withColumn("component", coalesce(col("component"), col("component2")))
  .groupBy("component")
  .agg(
    first(when(col("GUID").isNotNull, col("GUID"))).alias("GUID"),
    first(when(col("WID").isNotNull, col("WID"))).alias("WID"),
    first(when(col("SUid").isNotNull, col("SUid"))).alias("SUid")
  )

// 5. 关联原表得到最终结果
val result_df = df1.join(
  component_info,
  (df1("GUID") === components("id") && components("component") === component_info("component")) ||
  (df1("WID") === components("wid_id") && components("component2") === component_info("component")),
  "left"
).select(
  df1("EventId"),
  component_info("GUID"),
  component_info("WID"),
  component_info("SUid")
).distinct()

result_df.show()

注意:使用GraphFrames需额外添加依赖,Maven依赖示例:

<dependency>
  <groupId>org.graphframes</groupId>
  <artifactId>graphframes_2.12</artifactId>
  <version>0.8.2-spark3.2-s_2.12</version>
</dependency>

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 03:15:08