Spark/Scala中如何将collect_set返回的null替换为空白数组[]
解决Spark左外连接聚合后生成含空结构体数组的问题
左外连接customerDS与accountDS后分组聚合,尝试用coalesce+collect_set生成accounts字段,但输出中仍出现[{null, null, null}],需要将其替换为空数组[]。
原代码片段
val customerAccountDS: Dataset[CustomerAccountOutput] = customerDS .join(accountDS, customerDS("customerId") === accountDS("customerId"), "left_outer") .groupBy(customerDS.col("customerId"), customerDS.col("forename"), customerDS.col("surname")) .agg( coalesce(collect_set(struct(accountDS.col("customerId"), accountDS.col("accountId"), accountDS.col("balance"))), array()).alias("accounts")
问题原因
左外连接时,无匹配账户的客户会生成全字段为null的结构体,collect_set会将这个空结构体收集为数组(而非返回null),因此coalesce无法触发替换逻辑——它只处理null值,不处理包含空结构体的非空数组。
解决方案
方法1:聚合前过滤空结构体
在生成结构体时,仅保留账户字段非null的记录,无匹配的客户会直接生成空数组:
val customerAccountDS: Dataset[CustomerAccountOutput] = customerDS .join(accountDS, customerDS("customerId") === accountDS("customerId"), "left_outer") .groupBy(customerDS("customerId"), customerDS("forename"), customerDS("surname")) .agg( collect_set( when(accountDS("customerId").isNotNull, struct(accountDS("customerId"), accountDS("accountId"), accountDS("balance")) ) ).alias("accounts") ) // 移除数组中的null元素 .withColumn("accounts", filter(col("accounts"), x => x.isNotNull))
方法2:聚合后清理空数组
先完成聚合,再过滤数组中的全null结构体,若清理后数组为空则返回空数组:
import org.apache.spark.sql.functions._ val customerAccountDS: Dataset[CustomerAccountOutput] = customerDS .join(accountDS, customerDS("customerId") === accountDS("customerId"), "left_outer") .groupBy(customerDS("customerId"), customerDS("forename"), customerDS("surname")) .agg( collect_set(struct(accountDS("customerId"), accountDS("accountId"), accountDS("balance"))).alias("accounts") ) // 过滤数组中所有字段均为null的结构体 .withColumn("accounts", filter(col("accounts"), x => x("customerId").isNotNull || x("accountId").isNotNull || x("balance").isNotNull ) )
效果验证
修改后,无匹配账户的客户对应的accounts字段将显示为[],符合预期。
内容的提问来源于stack exchange,提问作者sarabanv
相关产品推荐
相关产品推荐

