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实现方案
核心步骤
- 定义JSON数组的Schema:明确
dsud_answers中JSON数组的结构,用于from_json转换。 - 字符串转JSON数组:使用
from_json将字符串列转换为Spark的数组类型。 - 过滤空数组:排除
dsud_answers为空数组的记录。 - 展开数组:用
explode将数组拆分为多行,每个元素对应一条记录。 - 数据转换与去重:处理日期格式、去重,最后重组所需列。
完整代码
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
相关产品推荐
相关产品推荐

