Spark DataFrame多ID关联参考表拼接字段及非空处理需求
Spark实现关联参考表后多字段条件拼接为单列
整体思路
我来一步步拆解这个需求的实现方式:首先得把主表的id1/id2/id3和参考表关联,拿到对应的描述和父类信息;然后根据id2和id3是否都非空的条件,决定是按顺序拼接多个值(末尾加/M)还是只用id1对应的内容。这里关键要保证拼接顺序和原ID字段的顺序一致,所以我们会用posexplode来保留ID的位置索引。
步骤1:导入必要的Spark函数
先把Spark SQL的核心函数导入进来,方便后续数据处理:
import org.apache.spark.sql.functions._
步骤2:定义原始数据集
先把你提供的主表和参考表准备好:
// 主数据集 val df = sc.parallelize(Seq( ("a", 1,2,3), ("b", 4,6,5) )).toDF("value", "id1", "id2", "id3") // 参考数据集 val ref = sc.parallelize(Seq( (1,"apple","fruit"), (2,"banana","fruit"), (3,"cat","animal"), (4,"dog","animal"), (5,"elephant","animal"), (6,"Flight","object") )).toDF("id", "descr", "parent")
步骤3:保留ID位置并关联参考表
为了保证拼接顺序严格遵循id1→id2→id3,我们用posexplode把ID数组拆分成带位置索引的行,再和参考表关联获取对应的描述信息:
// 将主表的三个ID转为数组并拆分,保留位置索引 val dfWithPos = df.withColumn("id_pos", posexplode(array(col("id1"), col("id2"), col("id3")))) .select( col("value"), col("id_pos.pos").alias("pos"), // 位置索引,确保拼接顺序不混乱 col("id_pos.col").alias("id") ) // 和参考表关联,拿到每个ID对应的descr和parent val dfJoined = dfWithPos.join(ref, Seq("id"), "left")
步骤4:按主表标识聚合值列表
按value分组,收集每个标识对应的descr和parent列表,同时过滤掉空值(避免ID为空时引入无效内容):
val dfAgg = dfJoined.groupBy("value") .agg( collect_list(when(col("id").isNotNull, col("descr"))).alias("desc_list"), collect_list(when(col("id").isNotNull, col("parent"))).alias("parent_list") )
步骤5:添加条件判断生成最终结果
回到原主表,先判断id2和id3是否都不为空,再和聚合后的数据集关联,根据条件生成最终的拼接结果:
// 添加条件字段:标记id2和id3是否均非空 val dfWithCheck = df.withColumn("is_all_valid", col("id2").isNotNull && col("id3").isNotNull) // 关联聚合结果,生成最终的desc和parent字段 val finalDF = dfWithCheck.join(dfAgg, Seq("value"), "inner") .withColumn("desc", when( col("is_all_valid"), concat_ws("+", col("desc_list"), lit("/M")) // 符合条件则拼接并加/M ).otherwise(col("desc_list").getItem(0)) // 不符合则只用id1的内容 ) .withColumn("parent", when( col("is_all_valid"), concat_ws("+", col("parent_list"), lit("/M")) ).otherwise(col("parent_list").getItem(0)) ) .select("desc", "parent") // 查看最终结果 finalDF.show(false)
运行结果
执行完上述代码后,会得到你期望的输出:
+-----------------------+--------------------------+ |desc |parent | +-----------------------+--------------------------+ |apple+banana+cat/M |fruit+fruit+animal/M | |dog+Flight+elephant/M |animal+object+animal/M | +-----------------------+--------------------------+
额外说明
- 如果主表中存在
id2或id3为空的行(比如("c",7,null,8)),最终结果会只显示id1对应的内容,不会触发拼接逻辑。 - 用
posexplode而不是普通explode,是为了确保拼接顺序和原ID字段顺序完全一致,避免数据错乱。
内容的提问来源于stack exchange,提问作者John Humanyun
相关产品推荐
相关产品推荐

