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

如何修复PySpark中修改数组内复杂结构体的代码问题

问题场景

给定一个DataFrame,每行是一个结构体,其中包含answers字段,该字段是一个由多个answer结构体组成的数组,每个answer结构体包含多个字段。以下代码用于处理数组中的每个answer,检查其render字段并进行处理(注:此代码运行在AWS Glue 3.0 Notebook中,除Spark上下文创建外,适用于所有PySpark >=3.1版本):

%glue_version 3.0

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
import pyspark.sql.functions as F
import pyspark.sql.types as T

sc = SparkContext.getOrCreate()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)

test_schema = T.StructType([
    T.StructField('question_id', T.StringType(), True),
    T.StructField('subject', T.StringType(), True),
    T.StructField('answers', T.ArrayType(
        T.StructType([
            T.StructField('render', T.StringType(), True),
            T.StructField('encoding', T.StringType(), True),
            T.StructField('misc_info', T.StructType([
                T.StructField('test', T.StringType(), True)
            ]), True)
        ]), True), True)
    ])

json_df = spark.createDataFrame(data=[
    [1, "maths", [("[tex]a1[/tex]", "text", ("x",)),("a2", "text", ("y",))]],
    [2, "bio", [("b1", "text", ("z",)),("<p>b2</p>", "text", ("q",))]],
    [3, "physics", None]
], schema=test_schema)
json_df.show(truncate=False)
json_df.printSchema()

运行结果如下:

+-----------+-------+---------------------------------------------+
|question_id|subject|answers                                      |
+-----------+-------+---------------------------------------------+
|1          |maths  |[{[tex]a1[/tex], text, {x}}, {a2, text, {y}}]|
|2          |bio    |[{b1, text, {z}}, {<p>b2</p>, text, {q}}]    |
|3          |physics|null                                         |
+-----------+-------+---------------------------------------------+

root
 |-- question_id: string (nullable = true)
 |-- subject: string (nullable = true)
 |-- answers: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- render: string (nullable = true)
 |    |    |-- encoding: string (nullable = true)
 |    |    |-- misc_info: struct (nullable = true)
 |    |    |    |-- test: string (nullable = true)

文本处理方法如下:

import re

@F.udf(returnType=T.StringType())
def determine_encoding_udf(render):
    if render:
        if "[tex]" in render:
            return "tex"
        
        match = re.search(r"<[^>]*>", render)
        if match:
            return "html"
    else:
        return "null render"
        
    return "text"

应用转换操作:

def normalize_answer(answer):
    return answer.withField(
        "processed_render_input",
        answer.getField("render")
    ).withField(
        "encoding",
        determine_encoding_udf(answer.getField("render"))
    )

json_mod_df = json_df.withColumn(
    "answers", 
    F.transform("answers", normalize_answer)
)

json_mod_df.show(truncate=False)

转换结果:

+-----------+-------+------------------------------------------------------------------------------+
|question_id|subject|answers                                                                       |
+-----------+-------+------------------------------------------------------------------------------+
|1          |maths  |[{[tex]a1[/tex], null render, {x}, [tex]a1[/tex]}, {a2, null render, {y}, a2}]|
|2          |bio    |[{b1, null render, {z}, b1}, {<p>b2</p>, null render, {q}, <p>b2</p>}]        |
|3          |physics|null                                                                          |
+-----------+-------+------------------------------------------------------------------------------+

process_text逻辑复杂,无法在transform的lambda表达式中实现。

问题描述

在处理更大的answers数据集(约5万行)时,encoding字段的内容对应其他行数组中的answer结构体的render值,甚至有时缺失;而在上述测试数据集中,所有encoding字段都显示"null render",但processed_render_input字段的值是对应answer的正确render内容。希望使用较新的transform和withField函数解决该问题,而非使用explode的传统处理方式。

临时解决方案

这并非理想的解决方案,但可以实现需求,在此分享给遇到相同问题的用户。若要获得正式解决方案,需要展示如何通过转换操作实现需求,本方案使用了explode函数,适用于复杂数据结构。

内容的提问来源于stack exchange,提问作者LaserJesus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 00:25:28