You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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.

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 the forminfo array.
  • struct_element -> struct_element.code = 'APPROVED': For each struct in the array, checks if its code field equals "APPROVED". Only matching structs are kept in the new array.
  • The resulting column forminfo_approved will have the exact same type as forminfo: 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 forminfo value, 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.27 17:07:48