Skip to content
sparksql-json-processing logo

sparksql-json-processing

SparkSQL JSON 处理完全指南。当用户需要处理 **Spark SQL 中的 JSON 字符串列** 时务必触发此技能!无论是简单字段提取还是复杂嵌套数据转换,此技能都能提供最优解。 务必在以下场景触发: - 提取 JSON 字段:使用 get_json_object 或 from_json - JSON 结构转换:重命名 key、增删字段、重组嵌套结构 - 数组处理:过滤(filter)、变换(transform)、展开 JSON 数组 - JSON Schema 推断:使用 schema_of_json - 完整 ETL 流程:解析 → 处理 → 重组 → 序列化 特别...

kpretty/skills0installs0starsOther

SKILL.md

Full skill instructions

SparkSQL JSON 处理技能

核心函数速查

场景函数
按 JSONPath 提取单值get_json_object(json, '$.field')
JSON → 结构化类型from_json(json_str, schema)
结构化类型 → JSONto_json(struct_expr)
按名称构造 Structnamed_struct('key', value, ...)
过滤数组元素filter(array, x -> condition)
转换数组元素transform(array, x -> expression)
推断 Schemaschema_of_json(json_literal)

核心链路

处理 JSON 数据最常用的完整链路:

解析 → 数组处理 → 重组 → 输出
from_json → filter/​transform → named_struct → to_json

为什么这个链路如此重要:JSON 字符串在 Spark 中是 TEXT 类型,无法直接用 SQL 操作。必须先 from_json 解析为 STRUCT/​ARRAY 等结构化类型,才能用点号取字段、用高阶函数处理数组。

基础模式

1. 提取嵌套字段

-- get_json_object:适合提取单个字段
SELECT 
  get_json_object(raw_json, '$.user.name') AS user_name,
  get_json_object(raw_json, '$.items[0].sku') AS first_sku
FROM table;

-- 注意:返回值始终是 STRING,数值运算需 CAST
SELECT CAST(get_json_object(raw_json, '$.amount') AS DOUBLE) * 0.9 AS discounted
FROM table;

2. 解析为结构化类型

-- from_json:将 JSON 字符串解析为 STRUCT,一次解析,多次使用
SELECT 
  parsed.id,
  parsed.user.name,    -- 点号访问嵌套字段
  parsed.items[0].sku  -- 下标访问数组元素
FROM (
  SELECT from_json(raw_json, 'STRUCT<id: INT, user: STRUCT<name: STRING>, items: ARRAY<STRUCT<sku: STRING>>>') AS parsed
  FROM table
);

3. JSON Key 重命名

-- 解析 → named_struct 重命名 → to_json 序列化
SELECT to_json(
  named_struct(
    'new_name', parsed.old_name,
    'another',   parsed.field_b
  )
) AS renamed_json
FROM (
  SELECT from_json(raw_json, 'STRUCT<old_name: STRING, field_b: INT>') AS parsed
  FROM table
);

4. 过滤和转换数组

-- filter:按条件过滤数组元素
filter(items, item -> item.qty > 1)

-- transform:对每个元素做映射
transform(items, item -> item.sku)

-- 组合:先过滤再转换
transform(filter(items, item -> item.price >= 100), item -> item.sku)

5. 完整链路示例

SELECT
  to_json(named_struct(
    'order_id',       parsed.order_id,
    'high_value_items', transform(
                          filter(parsed.items, i -> i.price >= 1000),
                          i -> named_struct(
                            'sku',      i.sku,
                            'name',     i.name,
                            'subtotal', ROUND(i.price * i.qty, 2)
                          )
                        )
  )) AS result
FROM (
  SELECT from_json(raw_json, '
    STRUCT<order_id: STRING, items: ARRAY<STRUCT<sku: STRING, name: STRING, price: DOUBLE, qty: INT>>>
  ') AS parsed
  FROM table
);

Schema DDL 语法

基本类型:  INT, BIGINT, DOUBLE, FLOAT, STRING, BOOLEAN, DATE, TIMESTAMP
结构体:    STRUCT<field1: TYPE1, field2: TYPE2>
数组:      ARRAY<TYPE>
Map:       MAP<KEY_TYPE, VALUE_TYPE>

最佳实践

  1. 优先用 from_json 一次性解析:对同一列只调用一次,避免重复解析同一 JSON 字符串

  2. filter 在前,transform 在后:先过滤缩小规模,再做映射转换

  3. named_struct + to_json 控制输出:需要输出 JSON 时,这是 Key 重命名和字段增删的核心工具

  4. schema_of_json 仅用于探索:生产环境应写死 Schema,避免推断偏差

  5. 注意 NULL 处理:JSON 中缺失字段经 from_json 后为 NULL,用 COALESCE 兜底

详细参考

遇到复杂场景时,查阅 references/​spark_json_guide.md 获取完整函数定义、使用案例和更多示例。