如何过滤Spark DataFrame数组内struct的hardboundFlag并保留全量列?
Spark DataFrame数组嵌套结构体过滤问题
问题详情
DataFrame结构
|-- student id: string |-- books: array | |-- element: struct | | |-- student id: string | | |-- book id: string | | |-- pages: string | | |-- hardboundFlag: boolean
原始数据
| St Id | Books |
|---|---|
| 1 | [1,1,200, true],[1,4,500, true] |
| 2 | [2,2,300, false] |
| 3 | [3,3,400, true],[3,5,100, false] |
需求
过滤books数组中hardboundFlag = true的结构体元素,同时排除过滤后书籍数组为空的学生记录,预期结果:
| St Id | Books |
|---|---|
| 1 | [1,1,200, true],[1,4,500, true] |
| 3 | [3,3,400, true] |
尝试使用explode方法后,St Id=1的记录仅返回单本书籍,不符合预期,求正确解决方法。
解决方法
直接使用Spark的数组高阶函数过滤数组内的元素,无需拆分原行结构,具体实现如下:
PySpark 代码
from pyspark.sql import functions as F # 过滤数组中符合条件的书籍,再剔除空数组记录 filtered_df = ( df.withColumn( "filtered_books", F.filter("books", lambda book: book["hardboundFlag"] == True) ) .filter(F.size("filtered_books") > 0) .drop("books") .withColumnRenamed("filtered_books", "books") ) # 查看结果 filtered_df.show(truncate=False)
Scala 代码
import org.apache.spark.sql.functions._ val filteredDF = df .withColumn( "filtered_books", filter(col("books"), book => book.getAs[Boolean]("hardboundFlag")) ) .filter(size(col("filtered_books")) > 0) .drop("books") .withColumnRenamed("filtered_books", "books") // 查看结果 filteredDF.show(false)
关键说明
F.filter(PySpark)/filter(Scala):直接对数组列进行过滤,保留所有满足hardboundFlag = true的结构体元素,不会拆分原有的学生行记录。size(col("filtered_books")) > 0:过滤掉过滤后书籍数组为空的学生(比如St Id=2的记录)。- 最后通过
drop和withColumnRenamed恢复原列名结构。
这种方法既能保留符合条件的所有书籍,又能保证学生记录的完整性,完全匹配预期结果。
内容的提问来源于stack exchange,提问作者anoop
相关产品推荐
相关产品推荐

