PySpark中高效读取拼接JSON文件的最优方案
问题描述
我有一个每行都是独立JSON对象的拼接JSON文件,示例内容如下:
{"id":"test1", "results": [{"property1": "sample1"},{"property2": "sample2"}]} {"id":"test2", "results": [{"property1": "sample3"},{"property2": "sample4"}]}
使用spark.read.json(filepath)读取时,只能获取第一个对象:
+-----+--------------------+ | id| results| +-----+--------------------+ |test1|[{sample1, null},...| +-----+--------------------+
期望得到的结构化结果是:
+-----+---------+---------+ |id |property1|property2| +-----+---------+---------+ |test1|sample1 |sample2 | |test2|sample3 |sample4 | +-----+---------+---------+
当前的处理方式是逐行读取文本、解析JSON后合并DataFrame,代码如下:
df = (spark.read.text(data[self.files])) dataCollect = df.collect() i = 0 for row in dataCollect: df_row = flatten_json(spark.read.json(spark.sparkContext.parallelize(row))) if i == 0: df_all = df_row else: df_all = df_row.unionByName(df_all, allowMissingColumns = True) i = i + 1
其中flatten_json是自动扁平化JSON的辅助函数,但这种方法效率较低,希望得到更优的处理方式。
优化解决方案
1. 正确读取每行JSON
Spark原生支持读取每行一个独立JSON对象的文件,只需确保读取配置正确(默认multiLine=False,即按行解析),无需手动逐行处理:
raw_df = spark.read.json(filepath, multiLine=False)
这一步就能直接得到包含所有行的DataFrame,解决仅读取第一行的问题。
2. 扁平化results数组并提取目标属性
根据results数组的结构,有两种高效处理方式:
方式一:通用场景(数组元素不固定顺序/数量)
通过explode展开数组,再用pivot聚合属性:
from pyspark.sql import functions as F # 展开results数组为单行记录 exploded_df = raw_df.select("id", F.explode("results").alias("result")) # 提取map的键值对,按id聚合得到扁平化列 flattened_df = exploded_df.groupBy("id")\ .pivot(F.map_keys("result")[0])\ .agg(F.first(F.map_values("result")[0]))
方式二:固定结构场景(已知数组含2个固定属性对象)
如果results数组始终是[property1对象, property2对象]的固定结构,直接索引提取更高效:
from pyspark.sql import functions as F flattened_df = raw_df.select( "id", F.col("results")[0]["property1"].alias("property1"), F.col("results")[1]["property2"].alias("property2") )
3. 替代自定义flatten_json(通用扁平化)
如果需要处理更复杂的嵌套JSON,可使用Spark内置函数实现通用扁平化,避免自定义函数的性能开销:
from pyspark.sql import functions as F # 合并results数组中的所有map为单个map merged_map_df = raw_df.select( "id", F.aggregate( "results", F.create_map().cast("map<string,string>"), lambda acc, x: F.map_concat(acc, x) ).alias("merged_properties") ) # 将map展开为对应的列 final_df = merged_map_df.select( "id", *[F.col("merged_properties")[key].alias(key) for key in ["property1", "property2"]] )
为什么当前方法效率低下?
collect()会将全量数据拉取到Driver节点,大数据量下极易引发内存溢出,完全丧失Spark的分布式处理优势。- 循环执行
unionByName会反复创建新的DataFrame,带来大量性能开销,数据量越大,效率越低。
内容的提问来源于stack exchange,提问作者Valouf
相关产品推荐
相关产品推荐

