Spark DataFrame 实战:清洗嵌套 JSONL 日志,轻松过滤订单数据
关键词Spark、DataFrame、JSONL、嵌套结构、struct、array、过滤引言在实际数据工程中我们经常要处理来自日志系统、API 接口的JSONL 文件每行一个 JSON 对象这些数据往往有多层嵌套例如订单中包含用户对象和商品数组。使用 Spark 的DataFrame API我们可以像操作 SQL 表一样用点号下钻 struct、用函数判断数组元素轻松完成复杂的过滤和字段拍平。本文将基于一个订单日志样例演示如何用 DataFrame 过滤出“收货城市为上海、用户年龄≥25、且订单包含 SKU-A”的记录并输出为 JSON。我们会展示代码、运行结果并深入讲解 DataFrame 与 RDD 的区别以及嵌套数据的处理技巧。实验环境同博客一Spark 4.2.0 Standalone 集群数据目录/workspace/data/jsonl/。样例数据orders.jsonl我们生成了 20 条订单数据固定随机种子可复现。每行 JSON 结构如下缩进展示{order_id:ORD-1001,user:{name:陈静,age:33,address:{city:Shanghai,district:区1,country:CN}},items:[{sku:SKU-A,price:32.01,qty:3},{sku:SKU-A,price:127.14,qty:1},{sku:SKU-C,price:375.52,qty:1}],total:598.69,status:paid}Spark 自动推断出的 Schema 如下通过df.printSchema()获得root |-- order_id: string (nullable true) |-- user: struct (nullable true) | |-- name: string (nullable true) | |-- age: long (nullable true) | |-- address: struct (nullable true) | | |-- city: string (nullable true) | | |-- country: string (nullable true) | | |-- district: string (nullable true) |-- items: array (nullable true) | |-- element: struct (containsNull true) | | |-- price: double (nullable true) | | |-- qty: long (nullable true) | | |-- sku: string (nullable true) |-- total: double (nullable true) |-- status: string (nullable true)完整代码filter_orders.py实验二多层嵌套 JSONL 过滤—— DataFrame API 版 用法 spark-submit --master master-url filter_orders.py 输入目录 输出目录 importsysfrompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimportarray_contains,coldefmain():input_path,output_pathsys.argv[1],sys.argv[2]sparkSparkSession.builder.appName(FilterOrders).getOrCreate()dfspark.read.json(input_path)# 自动推断嵌套 schemadf.printSchema()result(df.filter((col(user.address.city)Shanghai)# 嵌套 struct 下钻(col(user.age)25)# 嵌套数值过滤array_contains(col(items.sku),SKU-A)# 订单里含 SKU-A).select(order_id,col(user.name).alias(user_name),# 嵌套字段拍平改名col(user.address.city).alias(city),total,status,).cache()# 下面 count/write/show 三个 action缓存避免重复计算)print(f 过滤后行数:{result.count()})result.write.json(output_path)result.show(truncateFalse)spark.stop()if__name____main__:main()运行命令# 生成样例数据如需dockerexecspark python3 /workspace/scripts/gen_orders.py# 清除旧输出dockerexecsparkrm-rf/workspace/data/out/jsonl_filtered# 提交作业dockerexecspark /opt/spark/bin/spark-submit\--masterspark://localhost:7077\/workspace/scripts/filter_orders.py\/workspace/data/jsonl /workspace/data/out/jsonl_filtered真实输出控制台打印 过滤后行数: 2 -------------------------------------- |order_id|user_name|city |total |status| -------------------------------------- |ORD-1001|陈静 |Shanghai|598.69 |paid | |ORD-1006|张伟 |Shanghai|1012.23|unpaid| --------------------------------------输出目录下生成了带随机后缀的 JSON 文件这是 DataFrame 写出的正常行为$lsdata/out/jsonl_filtered/ _SUCCESS part-00000-5d605f45-909a-4419-8512-b2aafbea77f6-c000.json $catdata/out/jsonl_filtered/part-*.json{order_id:ORD-1001,user_name:陈静,city:Shanghai,total:598.69,status:paid}{order_id:ORD-1006,user_name:张伟,city:Shanghai,total:1012.23,status:unpaid}用 Python 独立复算原始 JSONL结果一致验证正确性。原理解读DataFrame 的威力DataFrame 与 RDD 的本质区别RDD 是无结构的数据集合JVM 对象而 DataFrame 带有Schema表结构——每列有名字和类型。Spark 的Catalyst 优化器能利用这个结构生成高效的物理执行计划例如列裁剪、谓词下推因此 DataFrame 通常比手写 RDD 代码更快。处理嵌套结构的“三板斧”访问 struct 字段用点号如col(user.address.city)。判断数组是否包含某元素用array_contains(col(items.sku), SKU-A)其中items.sku会自动将数组中的 struct 投影为 sku 字符串数组。拍平嵌套字段用col(user.name).alias(user_name)重命名将嵌套字段提升为顶级列。关于 Schema 推断spark.read.json会扫描全部数据来推断类型数据量大时可能较慢。生产环境建议显式提供 Schema使用StructType既加速读取又避免类型猜测错误。输出文件命名DataFrame 的write.json生成的文件名带有随机 UUID 后缀与 RDD 的固定part-00000不同但读取时统一读目录即可。常见问题multiline选项默认multilinefalse适合 JSONL如果 JSON 对象跨多行需设置.option(multiline, true)。null 值处理访问 null struct 的字段不会报错返回 null过滤条件会自然排除执行计划中会有isnotnull检查。输出覆盖write.json也不支持覆盖需要先删除目录或使用.mode(overwrite)。总结DataFrame API 让我们以声明式的方式处理复杂嵌套数据代码简洁、执行高效。通过本实验你学会了如何读取 JSONL、访问 struct 和 array、组合过滤条件并输出清洗后的数据。这是日常数据清洗的标配技能。下一步我们将进入更“硬核”的场景数据以 Parquet 列式存储其中一列是 JSON 字符串——如何解析并过滤请看第三篇博客。