Spark读取多个Parquet文件后如何转为字典的字典或字典列表?
将Spark Parquet文件转换为字典的字典/列表
你当前的代码只是将每个Parquet文件读取为Spark DataFrame并存储,没有把DataFrame转换为目标的字典结构。以下是两种符合需求的实现方案:
方案1:生成字典的列表(所有行数据转为字典)
该方案会把所有文件的每一行数据转换成一个字典,最终合并为一个大列表:
import os import glob path = "你的文件路径" parquet_files = glob.glob(os.path.join(path, '*.parquet')) dict_list = [] # 仅处理前5个文件(和你原代码逻辑一致) for file in parquet_files[:5]: df = spark.read.parquet(file) # 将每行数据转为字典并添加到列表 # collect()会把分布式数据拉到Driver端,大数据量需注意内存压力 for row in df.collect(): dict_list.append(row.asDict())
也可以通过Pandas中转简化写法:
for file in parquet_files[:5]: df = spark.read.parquet(file) # 转Pandas后直接生成每行字典的列表,再合并到总列表 pandas_df = df.toPandas() dict_list.extend(pandas_df.to_dict('records'))
方案2:生成字典的字典(按文件分组存储列数据)
该方案会以文件索引为键,每个文件对应的列名-列数据字典为值,最终组成嵌套字典:
import os import glob path = "你的文件路径" parquet_files = glob.glob(os.path.join(path, '*.parquet')) dict_of_dicts = {} for idx, file in enumerate(parquet_files[:5]): df = spark.read.parquet(file) columns = df.columns # 构建列名到对应列数据列表的字典 column_data_dict = { col: df.select(col).rdd.flatMap(lambda x: x).collect() for col in columns } dict_of_dicts[idx] = column_data_dict
注意事项
collect()和toPandas()都会将分布式的Spark数据拉取到Driver节点的本地内存中,如果单个Parquet文件数据量极大,可能会导致内存溢出。若处理超大规模数据,建议优先考虑分布式处理逻辑,而非转为本地字典结构。
内容的提问来源于stack exchange,提问作者user18373817
相关产品推荐
相关产品推荐

