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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:32:43