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 流程:解析 → 处理 → 重组 → 序列化 特别...
SKILL.md
Full skill instructions
SparkSQL JSON 处理技能
核心函数速查
| 场景 | 函数 |
|---|---|
| 按 JSONPath 提取单值 | get_json_object(json, '$.field') |
| JSON → 结构化类型 | from_json(json_str, schema) |
| 结构化类型 → JSON | to_json(struct_expr) |
| 按名称构造 Struct | named_struct('key', value, ...) |
| 过滤数组元素 | filter(array, x -> condition) |
| 转换数组元素 | transform(array, x -> expression) |
| 推断 Schema | schema_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>
最佳实践
-
优先用
from_json一次性解析:对同一列只调用一次,避免重复解析同一 JSON 字符串 -
filter在前,transform在后:先过滤缩小规模,再做映射转换 -
named_struct + to_json控制输出:需要输出 JSON 时,这是 Key 重命名和字段增删的核心工具 -
schema_of_json仅用于探索:生产环境应写死 Schema,避免推断偏差 -
注意 NULL 处理:JSON 中缺失字段经
from_json后为 NULL,用COALESCE兜底
详细参考
遇到复杂场景时,查阅 references/spark_json_guide.md 获取完整函数定义、使用案例和更多示例。
