news 2026/9/7 4:50:34

Parquet 中的“数据套数据”:使用 `from_json` 解析 JSON 字符串列并过滤

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Parquet 中的“数据套数据”:使用 `from_json` 解析 JSON 字符串列并过滤

关键词: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选项(如PERMISSIVEFAILFAST)控制行为。


常见问题

  • JSONPath 大小写$.geo.city必须与 JSON 中的键名完全一致。
  • 数组元素访问:用$.items[0].sku取第一个商品的 SKU。
  • Parquet 无法直接查看:可以转存为 JSON/CSV 或使用df.show()预览。
  • 输出文件跳号:这是正常现象,因为 Spark 保留了上游分区编号,读取时用目录即可。

总结

本实验展示了如何在 Parquet 中处理“JSON 字符串列”这种高级场景。关键点:

  1. Parquet列式存储比 JSON 更高效,适合生产环境。
  2. 解析 JSON 字符串有两种途径:get_json_object即用即取,from_json结构化解析更便于后续操作。
  3. 显式提供 Schema 可以提升性能并避免类型错误。

通过这三个递进实验,你已经从 RDD 的“搬砖”学到 DataFrame 的“看图纸”,再到 Parquet+JSON 的“挖金子”,掌握了 Spark 数据处理的核心套路。接下来,你就可以在实际项目中灵活应用了。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/7 4:50:20

Android语音唤醒接入指南:讯飞AIKit集成实战与踩坑记录

简介&#xff1a;面向安卓开发者的科大讯飞AIKit语音唤醒功能完整工程资源&#xff0c;解决在Android Studio中从零接入AIKit SDK、配置唤醒词并调通语音唤醒的难题&#xff0c;适合初学语音交互或需要快速落地唤醒功能的移动端工程师。资源包共647个文件&#xff0c;压缩后45.…

作者头像 李华
网站建设 2026/9/7 4:45:20

FPGA低延迟车牌识别:自研脉动卷积阵列全流程实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/7 4:43:46

机器学习异常值检测与处理:从IQR到孤立森林的完整指南

我在实际做机器学习项目时&#xff0c;最怕的不是模型训练时间太长&#xff0c;而是数据清洗阶段漏掉了异常值。几个极端样本看起来不影响大局&#xff0c;却能让均值失效、让线性回归系数明显偏移、让聚类结果面目全非。异常值&#xff0c;也叫离群点&#xff0c;简单理解就是…

作者头像 李华
网站建设 2026/9/7 4:41:37

CSDN首页发布文章CSDN同步助手三维非凸空间中路径规划的搜索范式统一理论及其在无人机自主导航中的协同优化研究(Matlab代码实现)50 / 100摘要:会在推荐、列表等场景外露

&#x1f4a5;&#x1f4a5;&#x1f49e;&#x1f49e;欢迎来到本博客❤️❤️&#x1f4a5;&#x1f4a5; &#x1f3c6;博主优势&#xff1a;&#x1f31e;&#x1f31e;&#x1f31e;博客内容尽量做到思维缜密&#xff0c;逻辑清晰&#xff0c;为了方便读者。 &#x1f381…

作者头像 李华