如何用PySpark展平复杂JSON?Azure Synapse场景报错求助
问题解决:PySpark展平JSON时遇到"need struct type but got string"错误
问题背景
使用Azure Synapse将API数据转存为JSON文件后,通过PySpark Notebook展平复杂JSON以加载到SQL数据库,执行代码时触发如下错误:
AnalysisException: Can't extract value from reportData#191: need struct type but got string
错误原因
报错指向explode("reportData.current_residents")行,说明reportData字段被解析成了字符串类型,而非预期的嵌套结构。这通常是因为JSON文件中reportData对应的内容本身是嵌套的JSON字符串,而非直接的JSON对象,PySpark默认无法自动解析这种嵌套格式。
解决步骤
确认原始数据Schema
取消注释dfGetReportData.printSchema(),运行后查看response.result.reportData的类型,确认是否为string。解析字符串类型的reportData为Struct
使用from_json函数将字符串格式的reportData解析成结构化数据,可通过自动推断或手动定义Schema实现。重新执行展平逻辑
解析完成后,再对结构化的current_residents数组执行explode操作。
修改后的完整代码
from pyspark.sql.functions import explode, col, from_json, schema_of_json from pyspark.sql.types import ArrayType, StructType, StructField, StringType, IntegerType, DoubleType, DateType # 读取原始JSON文件 dfGetReportData = spark.read.option("multiline", "true").option("mode", "PERMISSIVE").json('xxxxxxxxxxxxx') # 自动推断reportData的JSON Schema(适合快速调试) sample_report_data = dfGetReportData.select(col("response.result.reportData").cast(StringType())) \ .filter(col("response.result.reportData").isNotNull()).first()[0] report_data_schema = schema_of_json(sample_report_data) # 解析字符串类型的reportData为结构化数据 dfReportData = dfGetReportData.select( "response.requestId", "response.code", from_json(col("response.result.reportData"), report_data_schema).alias("reportData") ) # 展平current_residents数组 dfRentRoll = dfReportData.select(explode(col("reportData.current_residents")).alias("current_residents")) # 提取所需字段 dfCurrent_Residents = dfRentRoll.select( "current_residents.property_name", "current_residents.property", "current_residents.lookup_code", "current_residents.bldg_unit", "current_residents.bldg", "current_residents.unit", "current_residents.floorplan_name", "current_residents.unit_type", "current_residents.space_option", "current_residents.sqft", "current_residents.unit_status", "current_residents.unit_occupancy_type", col("current_residents.unit_address").getItem(0).alias("unit_address"), "current_residents.bed", "current_residents.bath", "current_residents.resident_name", "current_residents.phone_number", "current_residents.email", "current_residents.occupants", "current_residents.original_lease_start", "current_residents.lease_id", "current_residents.lease_status", "current_residents.lease_occupancy_type", "current_residents.move_in_date", "current_residents.lease_start_date", "current_residents.lease_end_date", "current_residents.previous_lease_end_date", "current_residents.move_out_date", "current_residents.lease_term_name", "current_residents.lease_term", "current_residents.contract_length_months", "current_residents.occupied_length_months", "current_residents.market_rent" ) # 写入Parquet dfCurrent_Residents.write.mode('overwrite').parquet("abfss://entrata@steussynapseprod001.dfs.core.windows.net/Silver/entities/Reporting/RentRoll")
补充说明
- 如果自动推断的Schema类型识别不准确(比如日期、数字类型),可手动定义Schema,示例:
report_data_schema = StructType([ StructField("current_residents", ArrayType(StructType([ StructField("property_name", StringType()), StructField("property", StringType()), # 按实际字段类型补充其他字段 ]))) ]) - 确保读取JSON时
multiline参数设置正确,避免因JSON换行导致的解析错误。
内容的提问来源于stack exchange,提问作者ChisumTrail
相关产品推荐
相关产品推荐

