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

