如何在Databricks中使用Spark SQL解析复杂JSON数据?
在Databricks中用Spark SQL解析嵌套JSON的实用方案
针对你提供的嵌套JSON结构,以下是基于Spark SQL的分步解析方案,适配你熟悉SQL语法的习惯:
1. 加载JSON数据到临时视图
先把JSON数据导入Databricks临时视图,方便后续操作:
- 若为单条JSON字符串:
CREATE OR REPLACE TEMP VIEW raw_json_data AS SELECT '{你的完整JSON内容}' AS json_str;
- 若为DBFS存储的JSON文件:
CREATE OR REPLACE TEMP VIEW raw_json_data AS SELECT * FROM json.`/dbfs/path/to/your/json/file.json`;
2. 展开顶层数组updates
updates是数组类型,每条元素对应一个点位更新记录,用explode函数将数组拆分为多行:
CREATE OR REPLACE TEMP VIEW exploded_updates AS SELECT messageTimestamp, siteReference, updateCount, explode(updates) AS update_record FROM raw_json_data;
3. 提取嵌套结构体的基础字段
update_record是结构体类型,直接用.访问内部字段,同时处理部分记录缺失的可选字段(比如vehicleData、vehicleExecution):
CREATE OR REPLACE TEMP VIEW parsed_basic_fields AS SELECT messageTimestamp, siteReference, updateCount, update_record.eventTimestamp, update_record.spotReference, -- 提取vehicleSpot核心字段 update_record.spotInfo.vehicleSpot.familyType AS vs_family_type, update_record.spotInfo.vehicleSpot.type AS vs_spot_type, update_record.spotInfo.vehicleSpot.externalReference AS vs_external_ref, update_record.spotInfo.vehicleSpot.state AS vs_state, -- 处理可选的vehicleData字段,空值替换为空字符串 IFNULL(update_record.spotInfo.vehicleData.familyType, '') AS vd_family_type, IFNULL(update_record.spotInfo.vehicleData.type, '') AS vd_vehicle_type, IFNULL(update_record.spotInfo.vehicleData.externalReference, '') AS vd_external_ref, -- 处理可选的vehicleExecution字段 IFNULL(update_record.spotInfo.vehicleExecution.familyType, '') AS ve_family_type, IFNULL(update_record.spotInfo.vehicleExecution.type, '') AS ve_execution_type, -- 保留fields数组用于后续解析 update_record.spotInfo.vehicleSpot.fields AS vs_fields, IFNULL(update_record.spotInfo.vehicleData.fields, struct(array() AS c3FieldStrings, array() AS c3FieldIntegers, array() AS c3FieldDecimals, array() AS c3FieldDateTimes)) AS vd_fields, IFNULL(update_record.spotInfo.vehicleExecution.fields, struct(array() AS c3FieldStrings, array() AS c3FieldIntegers, array() AS c3FieldDecimals, array() AS c3FieldDateTimes)) AS ve_fields FROM exploded_updates;
4. 将fields中的键值对数组转成列
fields下的c3FieldStrings、c3FieldIntegers等是键值对数组,用inline展开数组后,通过pivot转成结构化列,更适合SQL查询:
示例:解析vehicleSpot的字符串字段
CREATE OR REPLACE TEMP VIEW vs_string_fields_pivoted AS SELECT spotReference, COALESCE(`Blocking Spot`, '') AS vs_blocking_spot, COALESCE(Disabled, '') AS vs_disabled, COALESCE(Spot, '') AS vs_spot_name, COALESCE(`External Reference`, '') AS vs_ext_ref, COALESCE(`Sub Warehouse`, '') AS vs_sub_warehouse, COALESCE(Site, '') AS vs_site, COALESCE(`Origin Restriction`, '') AS vs_origin_restriction, COALESCE(`Zone Path Info`, '') AS vs_zone_path_info FROM ( SELECT spotReference, name, COALESCE(value, '') AS value FROM parsed_basic_fields LATERAL VIEW inline(vs_fields.c3FieldStrings) AS field ) AS vs_strings PIVOT ( MAX(value) FOR name IN ( 'Blocking Spot', 'Disabled', 'Spot', 'External Reference', 'Sub Warehouse', 'Site', 'Origin Restriction', 'Zone Path Info' ) );
同理解析整数和日期字段
-- 解析vehicleSpot的整数字段 CREATE OR REPLACE TEMP VIEW vs_integer_fields_pivoted AS SELECT spotReference, COALESCE(Workflow, 0) AS vs_workflow_id, COALESCE(`Vehicle Zone Type`, 0) AS vs_vehicle_zone_type FROM ( SELECT spotReference, name, COALESCE(value, 0) AS value FROM parsed_basic_fields LATERAL VIEW inline(vs_fields.c3FieldIntegers) AS field ) AS vs_integers PIVOT ( MAX(value) FOR name IN ('Workflow', 'Vehicle Zone Type') ); -- 解析vehicleSpot的日期字段,转成TIMESTAMP类型 CREATE OR REPLACE TEMP VIEW vs_datetime_fields_pivoted AS SELECT spotReference, COALESCE(`Locked Timestamp`, CAST('1970-01-01' AS TIMESTAMP)) AS vs_locked_ts, COALESCE(`Outgoing Timestamp`, CAST('1970-01-01' AS TIMESTAMP)) AS vs_outgoing_ts, COALESCE(`Busy Timestamp`, CAST('1970-01-01' AS TIMESTAMP)) AS vs_busy_ts, COALESCE(`Confirmed Date`, CAST('1970-01-01' AS TIMESTAMP)) AS vs_confirmed_dt FROM ( SELECT spotReference, name, COALESCE(CAST(value AS TIMESTAMP), CAST('1970-01-01' AS TIMESTAMP)) AS value FROM parsed_basic_fields LATERAL VIEW inline(vs_fields.c3FieldDateTimes) AS field ) AS vs_datetimes PIVOT ( MAX(value) FOR name IN ('Locked Timestamp', 'Outgoing Timestamp', 'Busy Timestamp', 'Confirmed Date') );
5. 合并所有表得到最终结构化结果
把基础字段表和各个解析后的字段表通过spotReference关联,得到完整的结构化数据:
CREATE OR REPLACE TEMP VIEW final_parsed_data AS SELECT pb.* EXCEPT(vs_fields, vd_fields, ve_fields), vs_str.* EXCEPT(spotReference), vs_int.* EXCEPT(spotReference), vs_dt.* EXCEPT(spotReference) -- 如需加入vehicleData、vehicleExecution的解析字段,按同样方法关联即可 FROM parsed_basic_fields pb LEFT JOIN vs_string_fields_pivoted vs_str ON pb.spotReference = vs_str.spotReference LEFT JOIN vs_integer_fields_pivoted vs_int ON pb.spotReference = vs_int.spotReference LEFT JOIN vs_datetime_fields_pivoted vs_dt ON pb.spotReference = vs_dt.spotReference; -- 查询最终结果 SELECT * FROM final_parsed_data;
额外实用技巧
- 提前定义Schema解析:如果JSON结构固定,可通过
from_json结合预定义Schema一次性解析,效率更高:
CREATE OR REPLACE TEMP VIEW raw_json_data AS SELECT from_json( json_str, 'struct<messageTimestamp:string,siteReference:string,updateCount:int,updates:array<struct<eventTimestamp:string,spotReference:string,spotInfo:struct<vehicleSpot:struct<familyType:string,type:string,externalReference:string,state:string,displayState:string,fields:struct<c3FieldStrings:array<struct<externalReference:string,name:string,value:string>>,c3FieldIntegers:array<struct<externalReference:string,name:string,value:int>>,c3FieldDecimals:array<struct<externalReference:string,name:string,value:decimal>>,c3FieldDateTimes:array<struct<externalReference:string,name:string,value:string>>>>,vehicleData:struct<familyType:string,type:string,externalReference:string,state:string,displayState:string,fields:struct<c3FieldStrings:array<struct<externalReference:string,name:string,value:string>>,c3FieldIntegers:array<struct<externalReference:string,name:string,value:int>>,c3FieldDecimals:array<struct<externalReference:string,name:string,value:decimal>>,c3FieldDateTimes:array<struct<externalReference:string,name:string,value:string>>>>,vehicleExecution:struct<familyType:string,type:string,externalReference:string,state:string,displayState:string,fields:struct<c3FieldStrings:array<struct<externalReference:string,name:string,value:string>>,c3FieldIntegers:array<struct<externalReference:string,name:string,value:int>>,c3FieldDecimals:array<struct<externalReference:string,name:string,value:decimal>>,c3FieldDateTimes:array<struct<externalReference:string,name:string,value:string>>>>>>>' ) AS json_data FROM raw_json_data;
之后直接通过json_data.updates等路径访问字段即可。
- 空值处理:用
IFNULL、COALESCE处理可选字段的空值,避免查询报错。 - 类型转换:将日期字符串转为
TIMESTAMP类型,方便后续时间过滤、计算。
内容的提问来源于stack exchange,提问作者Nathan Sundararajan
相关产品推荐
相关产品推荐

