Spark SQL/Scala中提取并求和JSON字符串内Sent值的优化方案
Spark SQL/Scala 提取非标准JSON字符串中所有"Sent"值并求和
现有数据表结构如下(value字段为varchar类型,存储非标准JSON格式数据):
| id | value |
|---|---|
| 123 | {78kfcX={"Sent": 77, "Respond": 31, "NoResponse": 31}, 97Facz={"Sent": 45, "Respond": 31, "NoResponse": 31}} |
| 333 | {5mdzrZ={"Sent": 1, "Respond": 1, "NoResponset": 1}} |
需求:提取每个id对应的所有"Sent"值,若存在多个则求和,预期结果:
| id | sent |
|---|---|
| 123 | 122 |
| 333 | 1 |
以下是无需自定义UDF的简洁实现方案:
方案1:Spark SQL内置函数实现
核心思路:先将非标准JSON转换为标准格式,再解析为Map类型,展开后提取"Sent"字段求和。
SELECT id, SUM(CAST(json_extract(value_struct, '$.Sent') AS INT)) AS sent FROM ( SELECT id, -- 转换非标准JSON为标准格式:给顶级键添加双引号 regexp_replace( regexp_replace(value, '(\\w+)=', '\"$1\":'), '(?<=\\{)(\\w+)=', '\"$1\":' ) AS standard_json, -- 解析标准JSON为Map类型 from_json(standard_json, 'map<string, struct<Sent: int, Respond: int, NoResponse: int, NoResponset: int>>') AS json_map FROM your_table ) t -- 展开Map中的所有键值对 LATERAL VIEW explode(json_map) exploded AS key, value_struct GROUP BY id;
关键步骤说明
regexp_replace:匹配无引号的顶级键,添加双引号将非标准JSON转为Spark可解析的格式。from_json:将标准JSON字符串解析为Map类型,定义结构体兼容" NoResponse"和"NoResponset"的拼写差异。LATERAL VIEW explode:展开Map中的每个子对象,最后按id分组求和。
方案2:Scala DataFrame API实现
基于同样的逻辑,用DataFrame API实现更贴合Scala开发场景:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 输入DataFrame val data = Seq( (123, "{78kfcX={\"Sent\": 77, \"Respond\": 31, \"NoResponse\": 31}, 97Facz={\"Sent\": 45, \"Respond\": 31, \"NoResponse\": 31}}"), (333, "{5mdzrZ={\"Sent\": 1, \"Respond\": 1, \"NoResponset\": 1}}") ).toDF("id", "value") // 定义结构体Schema,兼容拼写差异 val structSchema = StructType(Seq( StructField("Sent", IntegerType), StructField("Respond", IntegerType), StructField("NoResponse", IntegerType, nullable = true), StructField("NoResponset", IntegerType, nullable = true) )) // 定义Map类型Schema val mapSchema = MapType(StringType, structSchema) val result = data // 转换为标准JSON .withColumn("standard_json", regexp_replace(regexp_replace(col("value"), "(\\w+)=", "\"$1\":"), "(?<=\\{)(\\w+)=", "\"$1\":\"")) // 解析JSON为Map .withColumn("json_map", from_json(col("standard_json"), mapSchema)) // 展开Map中的子对象 .selectExpr("id", "explode(json_map) as (key, value_struct)") // 按id分组求和Sent值 .groupBy("id") .agg(sum(col("value_struct.Sent")).alias("sent")) result.show()
优势
- 全程使用Spark内置函数,避免自定义UDF的繁琐逻辑。
- 预定义Schema提升解析效率,同时兼容字段拼写差异。
- 代码更简洁,符合Spark原生API的使用习惯。
内容的提问来源于stack exchange,提问作者tglacierboy
相关产品推荐
相关产品推荐

