PySpark实现结构体数组过滤:保留code为APPROVED的结构体生成新列
Hey Christie, I’ve run into this exact scenario before—let’s get your filtered array column sorted out properly. The key here is using Spark’s built-in functions (way more efficient than UDFs for this kind of task) to keep the type consistent, or fixing the UDF if you prefer that approach.
Method 1: Use Spark SQL filter with expr (Recommended)
Spark has a native filter function for arrays that lets you directly filter elements without writing a UDF. This will preserve the original array-of-structs type perfectly.
Here’s the code:
from pyspark.sql import functions as F # Apply the filter directly using expr df = df.withColumn( "forminfo_approved", F.expr("filter(forminfo, struct_element -> struct_element.code = 'APPROVED')") )
What’s happening here:
filter(forminfo, ...): Iterates over each element in theforminfoarray.struct_element -> struct_element.code = 'APPROVED': For each struct in the array, checks if itscodefield equals "APPROVED". Only matching structs are kept in the new array.- The resulting column
forminfo_approvedwill have the exact same type asforminfo:array<struct<id:string,code:string>>.
Method 2: Fix Your UDF (If You Prefer This Approach)
If you tried using a UDF before and it didn’t work, the most common issue is not specifying the correct return type (Spark needs to know exactly what schema the UDF returns). Here’s how to do it right:
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, StructType, StructField, StringType # Define the schema that matches your original forminfo column struct_schema = StructType([ StructField("id", StringType(), nullable=True), StructField("code", StringType(), nullable=True) ]) array_struct_type = ArrayType(struct_schema) # Define the UDF with null handling (important if your array can be null) def filter_approved_structs(arr): if arr is None: return None # Return null if input is null to avoid errors return [item for item in arr if item.get("code") == "APPROVED"] # Register the UDF with the correct return type filter_approved_udf = F.udf(filter_approved_structs, array_struct_type) # Apply the UDF to create the new column df = df.withColumn("forminfo_approved", filter_approved_udf(F.col("forminfo")))
Why this works:
- We explicitly define the return type as
ArrayType(struct_schema), which matches your original column’s type. - We handle null arrays to prevent runtime errors (if any rows have a null
forminfovalue, the UDF returns null instead of crashing).
Verify the Result
After running either method, check the column types to confirm they match:
print(df.dtypes)
You should see forminfo_approved listed with the same array<struct<id:string,code:string>> type as forminfo.
内容的提问来源于stack exchange,提问作者Christie

