如何在Spark DataFrame中过滤结构体数组内全为null的结构体
实现方案
你原有的思路存在偏差,不需要拿整个数组和空结构体做等值判断,直接用Spark内置的filter高阶函数遍历数组元素,过滤掉所有字段均为null的结构体即可,适配不同使用场景的实现如下:
PySpark实现
固定字段写法(已知结构体字段为id、type、name)
from pyspark.sql import functions as F df = df.withColumn( "brands", F.filter( F.col("brands"), lambda struct: ~(struct.id.isNull() & struct.type.isNull() & struct.name.isNull()) ) )
动态字段写法(自动适配结构体所有字段,无需手动枚举)
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StructType # 获取结构体的所有字段名 struct_schema = df.schema["brands"].dataType.elementType all_fields = struct_schema.names df = df.withColumn( "brands", F.filter( F.col("brands"), lambda struct: ~F.reduce( F.array([struct[field].isNull() for field in all_fields]), F.lit(True), lambda acc, is_null: acc & is_null ) ) )
Scala Spark实现
固定字段写法
import org.apache.spark.sql.functions._ val processedDf = df.withColumn( "brands", filter( col("brands"), struct => !(struct.getAs[String]("id").isNull && struct.getAs[String]("type").isNull && struct.getAs[String]("name").isNull) ) )
动态字段写法
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.{ArrayType, StructType} val structFields = df.schema("brands").dataType.asInstanceOf[ArrayType] .elementType.asInstanceOf[StructType].names val processedDf = df.withColumn( "brands", filter( col("brands"), struct => !structFields.map(f => struct.getAs[Any](f).isNull).reduce(_ && _) ) )
注意事项
- 以上实现基于Spark 3.0及以上版本,低于该版本可自行实现UDF完成数组过滤逻辑
- 直接判断结构体所有字段均为null的方式,比构造全空结构体做等值判断的兼容性和稳定性更高,不会出现因为结构体schema微小差异导致的判断失效问题
内容的提问来源于stack exchange,提问作者Meesam
相关产品推荐
相关产品推荐

