如何修复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
相关产品推荐
相关产品推荐

