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

如何在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;

额外实用技巧

  1. 提前定义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等路径访问字段即可。

  1. 空值处理:用IFNULL、COALESCE处理可选字段的空值,避免查询报错。
  2. 类型转换:将日期字符串转为TIMESTAMP类型,方便后续时间过滤、计算。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 01:42:05