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

PySpark实现Parquet字符串列转JSON类型并展开数据方案

PySpark 处理字符串格式JSON数组并展开的解决方案

问题背景

现有从SQL Server导出的Parquet文件,Schema如下:

root
 |-- user_uid: string (nullable = true)
 |-- user_email: string (nullable = true)
 |-- ud_id: integer (nullable = true)
 |-- ud_standard_workflow_id: integer (nullable = true)
 |-- ud_is_preview: boolean (nullable = true)
 |-- ud_is_completed: boolean (nullable = true)
 |-- ud_language: string (nullable = true)
 |-- ud_created_date: timestamp (nullable = true)
 |-- ud_modified_date: string (nullable = true)
 |-- ud_created_by_id: string (nullable = true)
 |-- dsud_id: integer (nullable = true)
 |-- dsud_user_data_id: integer (nullable = true)
 |-- dsud_dynamic_step_id: integer (nullable = true)
 |-- dsud_is_completed: boolean (nullable = true)
 |-- dsud_answers: string (nullable = true)

其中dsud_answers为字符串类型,实际存储JSON数组,示例:

[{"QuestionId":6406,"QuestionTitle":"Residency","Value":"1975"},{"QuestionId":6407,"QuestionTitle":"Citizentship","Value":"66664"}]

需要将该列转换为JSON类型并展开数组(每条数组元素生成一条记录),已用Pandas实现,但PySpark中报错:data type mismatch: Input schema "STRING" must be a struct, an array or a map,需提供PySpark实现方案。

附Pandas实现代码:

from pyspark.sql.types import StructField, IntegerType, TimestampType, BooleanType

def batch_function (df_answers, batch_id):
    

    df = df_answers.select("*").filter(df_bgx_answers.dsud_answers != '[]')
         .withColumn("ud_modified_date", to_timestamp(df_bgx_answers.ud_modified_date))
         .drop_duplicates()
         .toPandas()


    df_attributes = pd.DataFrame()
    df_final = pd.DataFrame()
    # Loop through the data to fill the dataframe
    for index in df.index:

        indexId = df.dsud_id[index]
        userDataId = df.dsud_user_data_id[index]
        dynamicStepId = df.dsud_dynamic_step_id[index]
        languageID = df.ud_language[index]
        createdDate = to_datetime(df.ud_created_date[index])
        createdBy = df.ud_created_by_id[index]
        modifiedDate = df.ud_modified_date[index]
        email = df.user_email[index]
        workflowId = df.ud_standard_workflow_id[index]
        uid = df.user_uid[index]
        completed = df.wrn_is_completed[index]
        agreed = df.wrn_is_agreed[index]
        flow_name = df.wf_name[index]

        row_json = json.loads(df.dsud_answers[index])
        normalized_row = pd.json_normalize(row_json)
        
        df_attributes = pd.concat([df_attributes, normalized_row], ignore_index=True) 
        df_attributes['dsud_user_data_id'] = userDataId
        df_attributes['dsud_id'] = indexId
        df_attributes['dsud_dynamic_step_id'] = dynamicStepId
        df_attributes['ud_language'] = languageID
        df_attributes['ud_created_date'] = createdDate
        df_attributes['ud_created_by_id'] = createdBy
        df_attributes['ud_modified_date'] = modifiedDate
        df_attributes['user_email'] = email
        df_attributes['ud_standard_workflow_id'] = workflowId
        df_attributes['user_uid'] = uid
        df_attributes['wrn_is_completed'] = completed
        df_attributes['wrn_is_agreed'] = agreed
        df_attributes['wf_name'] = flow_name

        df_attributes = df_attributes.reset_index(drop=True)
        df_final = pd.concat([df_final, df_attributes])

    df_answers = spark.createDataFrame(df_final)

    df_answers.write.mode("append").format("delta").saveAsTable("final_table")  

PySpark实现方案

核心步骤

  1. 定义JSON数组的Schema:明确dsud_answers中JSON数组的结构,用于from_json转换。
  2. 字符串转JSON数组:使用from_json将字符串列转换为Spark的数组类型。
  3. 过滤空数组:排除dsud_answers为空数组的记录。
  4. 展开数组:用explode将数组拆分为多行,每个元素对应一条记录。
  5. 数据转换与去重:处理日期格式、去重,最后重组所需列。

完整代码

from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, ArrayType

def batch_function(df_answers, batch_id):
    # 定义dsud_answers中JSON数组元素的结构Schema
    answer_item_schema = StructType([
        StructField("QuestionId", IntegerType(), nullable=True),
        StructField("QuestionTitle", StringType(), nullable=True),
        StructField("Value", StringType(), nullable=True)
    ])
    
    df_processed = (df_answers
        # 转换日期字段为timestamp类型
        .withColumn("ud_modified_date", F.to_timestamp(F.col("ud_modified_date")))
        # 将字符串格式的JSON数组转为Spark数组类型
        .withColumn("dsud_answers_array", F.from_json(F.col("dsud_answers"), ArrayType(answer_item_schema)))
        # 过滤空数组记录
        .filter(F.size(F.col("dsud_answers_array")) > 0)
        # 展开数组为多行,每个数组元素对应一条新记录
        .withColumn("answer_detail", F.explode(F.col("dsud_answers_array")))
        # 提取JSON元素中的字段
        .withColumn("QuestionId", F.col("answer_detail.QuestionId"))
        .withColumn("QuestionTitle", F.col("answer_detail.QuestionTitle"))
        .withColumn("Value", F.col("answer_detail.Value"))
        # 去重
        .dropDuplicates()
        # 选择最终需要的列(根据业务需求调整,若有Pandas中的wrn_is_completed等字段需自行添加)
        .select(
            "user_uid", "user_email", "ud_id", "ud_standard_workflow_id",
            "ud_is_preview", "ud_is_completed", "ud_language", "ud_created_date",
            "ud_modified_date", "ud_created_by_id", "dsud_id", "dsud_user_data_id",
            "dsud_dynamic_step_id", "dsud_is_completed", "QuestionId", "QuestionTitle", "Value"
        )
    )
    
    # 写入Delta表
    df_processed.write.mode("append").format("delta").saveAsTable("final_table")

关键说明

  • Schema匹配:from_json必须传入与JSON结构完全匹配的Schema,这里定义数组元素的结构体后,外层用ArrayType包裹,确保转换后是Spark可识别的数组类型,解决报错问题。
  • 空数组过滤:用size函数判断数组长度,比直接对比字符串'[]'更稳定,避免JSON格式细微差异导致的漏判。
  • 分布式处理:PySpark无需逐行循环,通过内置函数完成数组展开,完全利用分布式计算能力,性能远优于Pandas循环方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 07:12:17