Scala中Spark DataFrame ArrayType列移除空列表与空值的实现方法
Spark ArrayType列空值处理方案
分为内置函数实现和UDF实现两种方式:
方案1:内置函数实现(Spark 2.4+ 推荐,性能更优)
用到array_remove内置函数移除数组内的null元素,再通过size函数过滤处理后为空的数组行:
PySpark 示例
from pyspark.sql import functions as F # 替换为你自己的原始DataFrame变量名 processed_df = df.withColumn("xyz", F.array_remove(F.col("xyz"), None)) \ .filter(F.size(F.col("xyz")) > 0)
Scala 示例
import org.apache.spark.sql.functions.{array_remove, size, col} val processedDf = df.withColumn("xyz", array_remove(col("xyz"), null)) .filter(size(col("xyz")) > 0)
Spark SQL 示例
SELECT array_remove(xyz, null) AS xyz FROM 你的表名 WHERE size(array_remove(xyz, null)) > 0
方案2:UDF实现(兼容2.4以下版本Spark)
如果Spark版本较低没有array_remove函数,可以用自定义UDF实现相同逻辑:
PySpark 示例
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StringType # 自定义数组清洗逻辑 def clean_array(arr): if not arr: return [] return [elem for elem in arr if elem is not None] # 注册UDF,数组元素类型请根据实际场景调整 clean_arr_udf = F.udf(clean_array, ArrayType(StringType())) processed_df = df.withColumn("xyz", clean_arr_udf(F.col("xyz"))) \ .filter(F.size(F.col("xyz")) > 0)
执行逻辑说明
- 清洗步骤:将原数组内所有null元素移除,示例中
[a,,](即["a", null, null])会被处理为["a"],空数组[]处理后仍为[] - 过滤步骤:仅保留处理后数组长度大于0的行,过滤掉空数组行,得到预期结果。
内容的提问来源于stack exchange,提问作者Lumos
相关产品推荐
相关产品推荐

