关键词:Spark、Parquet、JSON 字符串、from_json、get_json_object、列式存储
引言
数据工程中常遇到一种场景:从 Kafka 或 HTTP 接口接收的数据,来不及解析全部字段,先以JSON 字符串的形式存入一列,后续再按需解析。这种“数据套数据”的模式在数据湖中很常见,尤其是 Parquet 列式存储文件。
本文将演示如何读取一个 Parquet 表,其中payload列是整段 JSON 字符串,我们需要拆开它,过滤出“城市为上海或北京、设备为 android、会话时长 ≥120 秒”的事件。我们会对比两种解析方式:get_json_object(取单个字段)和from_json(结构化解析),并解释 Parquet 的优势。
实验环境
依然使用 Spark 4.2.0 Standalone 集群。数据目录/workspace/data/parquet/events.parquet。
样例数据:events.parquet
我们使用脚本生成 20 条事件记录,保存为 Parquet 格式(Snappy 压缩)。表结构只有三列,payload是字符串:
root |-- event_id: string (nullable = true) |-- ts: string (nullable = true) |-- payload: string (nullable = true)而payload内部是一个结构化的 JSON,例如:
{"user_id":"u103","action":"click","device":{"os":"android","version":"11"},"geo":{"city":"Shanghai","country":"CN"},"session":{"duration_sec":474}}为了方便核对,同批数据还另存了一份events_reference.jsonl,因为 Parquet 是二进制文件,无法直接cat查看。
完整代码:filter_events.py
"""实验三:Parquet 中 JSON 字符串列的解析与过滤 用法: spark-submit --master <master-url> filter_events.py <输入 parquet 目录> <输出目录> """importsysfrompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportcol,from_json,get_json_objectfrompyspark.sql.typesimportIntegerType,StringType,StructField,StructType# from_json 需要显式给出 payload 的结构:字段带类型,解析后才能直接比较/聚合PAYLOAD_SCHEMA=StructType([StructField("user_id",StringType()),StructField("action",StringType()),StructField("device",StructType([StructField("os",StringType()),StructField("version",StringType()),])),StructField("geo",StructType([StructField("city",StringType()),StructField("country",StringType()),])),StructField("session",StructType([StructField("duration_sec",IntegerType()),])),])defmain():input_path,output_path=sys.argv[1],sys.argv[2]spark=SparkSession.builder.appName("FilterEvents").getOrCreate()df=spark.read.parquet(input_path)df.printSchema()# 方式一:get_json_object —— 顺手取一两个字段时用(返回字符串)df.select("event_id",get_json_object("payload","$.geo.city").alias("city_gjo"),).show(5,truncate=False)# 方式二:from_json + schema —— 结构化解析,字段带类型,适合过滤/聚合parsed=df.withColumn("p",from_json(col("payload"),PAYLOAD_SCHEMA)).select("event_id","ts",col("p.user_id").alias("user_id"),col("p.geo.city").alias("city"),col("p.device.os").alias("os"),col("p.session.duration_sec").alias("duration_sec"),)result=(parsed.filter(col("city").isin("Shanghai","Beijing")# 解析出的字段直接过滤&(col("os")=="android")&(col("duration_sec")>=120)# 类型已是 int,可直接比较).cache())print(f"===== 过滤后行数:{result.count()}=====")result.write.mode("overwrite").json(output_path)result.show(truncate=False)spark.stop()if__name__=="__main__":main()运行命令
# 生成样例(本地模式即可)dockerexecspark /opt/spark/bin/spark-submit--master'local[*]'/workspace/scripts/gen_events.py# 解析过滤并提交到集群dockerexecspark /opt/spark/bin/spark-submit\--masterspark://localhost:7077\/workspace/scripts/filter_events.py\/workspace/data/parquet/events.parquet\/workspace/data/out/events_filtered真实输出
首先打印原始 Schema(payload 是字符串):
root |-- event_id: string (nullable = true) |-- ts: string (nullable = true) |-- payload: string (nullable = true)用get_json_object取 city 字段预览:
+--------+--------+ |event_id|city_gjo| +--------+--------+ |EVT-2007|Shenzhen| |EVT-2008|Hangzhou| |EVT-2009|Shanghai| |EVT-2010|Hangzhou| |EVT-2017|Hangzhou| +--------+--------+最终过滤结果(3 条):
===== 过滤后行数: 3 ===== +--------+-------------------+-------+--------+-------+------------+ |event_id|ts |user_id|city |os |duration_sec| +--------+-------------------+-------+--------+-------+------------+ |EVT-2013|2026-09-06 10:13:00|u113 |Shanghai|android|337 | |EVT-2003|2026-09-06 10:03:00|u103 |Shanghai|android|474 | |EVT-2004|2026-09-06 10:04:00|u104 |Shanghai|android|594 | +--------+-------------------+-------+--------+-------+------------+输出目录为 JSON 格式,可以看到 3 个 part 文件(编号跳号,因为继承了上游 Parquet 的分区编号):
$lsdata/out/events_filtered/ _SUCCESS part-00000.snappy.parquet part-00003.snappy.parquet part-00007.snappy.parquet(注意:这里我们写的输出是 JSON,但示例中使用了write.json,实际文件是 JSON 格式,命名同样会跳号。)
原理解读:Parquet 与 JSON 解析
Parquet 列式存储:与 JSON 的行式存储不同,Parquet 按列存储数据,同列的数据连续存放,并应用压缩。这使得 Spark 可以只读取需要的列(例如只读payload而不读其他列),大大减少 I/O。因此,数据湖中的正式表常采用 Parquet。
两种 JSON 解析方式对比:
| 函数 | 适用场景 | 返回值类型 | 需提供 Schema |
|---|---|---|---|
get_json_object(col, path) | 偶尔取一两个字段 | 字符串(如需数值需 cast) | 否 |
from_json(col, schema) | 整体解析,后续多次使用 | 结构体,字段带原始类型 | 是(本实验的PAYLOAD_SCHEMA) |
本实验中,我们需要过滤duration_sec >= 120,这个字段是整数,如果用get_json_object得到字符串,还需要cast,而from_json直接按IntegerType解析,非常方便。
解析失败处理:默认情况下,如果某行 JSON 不符合 schema,from_json会返回 null。这可能导致过滤结果比预期少。生产环境中,可以使用mode选项(如PERMISSIVE、FAILFAST)控制行为。
常见问题
- JSONPath 大小写:
$.geo.city必须与 JSON 中的键名完全一致。 - 数组元素访问:用
$.items[0].sku取第一个商品的 SKU。 - Parquet 无法直接查看:可以转存为 JSON/CSV 或使用
df.show()预览。 - 输出文件跳号:这是正常现象,因为 Spark 保留了上游分区编号,读取时用目录即可。
总结
本实验展示了如何在 Parquet 中处理“JSON 字符串列”这种高级场景。关键点:
- Parquet列式存储比 JSON 更高效,适合生产环境。
- 解析 JSON 字符串有两种途径:
get_json_object即用即取,from_json结构化解析更便于后续操作。 - 显式提供 Schema 可以提升性能并避免类型错误。
通过这三个递进实验,你已经从 RDD 的“搬砖”学到 DataFrame 的“看图纸”,再到 Parquet+JSON 的“挖金子”,掌握了 Spark 数据处理的核心套路。接下来,你就可以在实际项目中灵活应用了。