You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 00:10:47