如何在HDInsight PySpark v2.x中读取大型JSON数组文件
处理HDInsight PySpark 2.x读取大型JSON数组文件的解决方案
我之前也碰到过类似的大体积JSON数组处理难题,结合PySpark 2.x的特性,给你几个可行的思路:
核心问题拆解
Spark默认的JSON读取器是基于JSON Lines格式(每行一个独立JSON对象)设计的,直接读整个JSON数组会解析失败,加上你的文件体积大、Schema不固定,提前定义Schema也不现实,所以核心方向是把JSON数组转成Spark能识别的格式,再利用自动Schema推断来处理。
方法1:用RDD预处理转成JSON Lines格式(最推荐)
这个方法利用RDD的分布式处理能力,先把大JSON数组拆成每行一个JSON对象,再转成DataFrame让Spark自动推断Schema,完全适配PySpark 2.x:
# 1. 读取大型JSON数组文件为RDD(分布式读取,避免单节点内存压力) json_rdd = sc.textFile("adl://your-datalake-path/large-file.json") # 2. 定义处理函数:去掉数组首尾的[],拆分每个JSON对象 def parse_json_array(line): # 处理整个文件是单行大数组的情况(大部分大JSON数组都是这样存储的) if line.strip().startswith('[') and line.strip().endswith(']'): # 去掉首尾括号后,按},{拆分每个对象 return line.strip()[1:-1].split('},{') # 处理文件已经有换行的情况 else: return [line.strip()] # 3. 扁平化RDD,得到每个元素是一个JSON对象的片段 flattened_rdd = json_rdd.flatMap(parse_json_array) # 4. 补全每个JSON对象的首尾括号(拆分后可能丢失) fixed_json_rdd = flattened_rdd.map(lambda x: '{' + x + '}' if not (x.startswith('{') and x.endswith('}')) else x ) # 5. 转成DataFrame,Spark会自动推断所有列的Schema(支持数百列) df = spark.read.json(fixed_json_rdd)
优化点:
如果文件特别大,自动推断Schema可能耗时,可以先采样1%-5%的数据提前推断Schema,再用该Schema读取全量数据:
# 采样部分数据推断Schema sample_rdd = json_rdd.sample(False, 0.01) # 采样1%的内容 processed_sample = sample_rdd.flatMap(parse_json_array).map(lambda x: '{' + x + '}' if not (x.startswith('{') and x.endswith('}')) else x) sample_df = spark.read.json(processed_sample) target_schema = sample_df.schema # 用预推断的Schema读取全量数据,提升速度 df = spark.read.schema(target_schema).json(fixed_json_rdd)
方法2:利用multiLine参数(仅限Spark 2.2+)
如果你的HDInsight PySpark是2.2及以上版本,可以尝试开启multiLine=True参数直接读取JSON数组,但注意:这个参数会把整个文件加载到一个Executor中,110GB的文件很容易触发OOM,只适合中小体积的JSON数组文件,代码示例:
# 仅适合Spark 2.2+,大文件慎用 df = spark.read.option("multiLine", "true").json("adl://your-datalake-path/large-file.json")
批量处理多个Schema不同的文件
因为你有大量Schema不同的文件,可以写一个循环遍历文件路径,逐个处理:
import os # 假设你的文件都在ADLS的某个目录下 file_paths = [f"adl://your-datalake-path/{file}" for file in os.listdir("/mnt/adls/your-directory")] for file_path in file_paths: # 复用上面的RDD处理逻辑 json_rdd = sc.textFile(file_path) flattened_rdd = json_rdd.flatMap(parse_json_array) fixed_json_rdd = flattened_rdd.map(lambda x: '{' + x + '}' if not (x.startswith('{') and x.endswith('}')) else x) df = spark.read.json(fixed_json_rdd) # 保存成Parquet格式(比JSON更适合后续分析,支持Schema自动存储) df.write.parquet(f"adl://your-datalake-path/processed/{os.path.basename(file_path)}.parquet")
集群资源配置建议
处理110GB的文件时,务必调整Spark的资源参数,避免内存溢出:
- 提交作业时设置:
--executor-memory 16G --executor-cores 4 --num-executors 10(根据你的集群节点数调整) - 开启Spark的内存溢出检查,调整
spark.driver.memory为8G以上
内容的提问来源于stack exchange,提问作者jatal
相关产品推荐
相关产品推荐

