Spark Scala中如何将一个DF字段添加到另一DF的嵌套struct字段中
Scala Spark 实现给嵌套Struct新增字段方案
核心逻辑是通过员工唯一标识左关联两个DataFrame,空值填充默认值后重构嵌套结构体,全程保留DF1的所有原始字段和记录。
import org.apache.spark.sql.functions._ import spark.implicits._ // 1. 打平DF2取关联所需字段,避免结构体字段名冲突 val df2Processed = df2.select( $"employee.id" as "join_key", $"employee.phone" as "source_phone" ) // 2. 左关联保留DF1全量数据,未匹配记录的phone填充默认值000 val joinedDF = df1.join(df2Processed, df1("employee.id") === df2Processed("join_key"), "left") .withColumn("final_phone", coalesce($"source_phone", lit("000"))) .drop("join_key", "source_phone") // 3. 自动读取DF1原有employee结构体的所有字段,重构包含phone的新结构体 val originEmpFields = df1.schema("employee").dataType .asInstanceOf[org.apache.spark.sql.types.StructType] .fieldNames val resultDF = joinedDF.withColumn( "employee", struct( originEmpFields.map(field => col(s"employee.$field")) :+ col("final_phone").as("phone"): _* ) ).drop("final_phone")
注意事项
- 关联默认使用
employee.id作为员工唯一主键,如果你的业务场景用其他字段做唯一标识(比如name+dept组合键),替换join条件里的关联字段即可 - 重构结构体时没有硬编码employee下的字段,会自动保留DF1中employee结构体下的所有原有字段(包括你省略的其他字段),不会出现字段丢失
coalesce函数会优先取DF2匹配到的phone值,只有关联不到记录(source_phone为null)时才会填充默认值"000",符合需求- 最终返回的
resultDF的schema中,employee结构体会新增phone: string字段,和预期结构一致
内容的提问来源于stack exchange,提问作者Nitish N Banakar
相关产品推荐
相关产品推荐

