Databricks中如何遍历Parquet文件路径并统计行列数至DataFrame
问题解决:统计所有Parquet文件的路径及行列数
原代码的问题
- 路径未批量收集:你只在循环里更新单个
path变量,最后构建字典时'path': path只会把最后一次循环的路径赋值给整列,导致所有行的路径都重复成最后一个文件的路径。 - 变量名错误:读取Parquet文件后赋值给了
df_parquet,但后面调用df.shape,这里df未定义,会触发报错。
修正后的Pandas代码
import os import pandas as pd # 初始化存储路径、行数、列数的列表 paths = [] row_counts = [] col_counts = [] # 遍历所有文件夹及文件 for dirpath, _, filenames in os.walk("/dbfs/mnt/sgi/report/"): # 筛选出.parquet结尾的文件 for filename in [f for f in filenames if f.endswith(".parquet")]: full_path = os.path.join(dirpath, filename) paths.append(full_path) # 读取Parquet文件并统计行列数 df_parquet = pd.read_parquet(full_path) row_counts.append(df_parquet.shape[0]) col_counts.append(df_parquet.shape[1]) # 构建DataFrame result_df = pd.DataFrame({ "path": paths, "count_of_rows": row_counts, "count_of_col": col_counts }) # 生成Excel文件(Databricks环境中保存到DBFS路径后可下载) result_df.to_excel("/dbfs/mnt/sgi/report/parquet_stats.xlsx", index=False) # 验证所有路径是否正确存储 print(result_df['path'].unique())
更高效的PySpark方案(适合大文件/多文件场景)
如果你的Parquet文件数量多、体积大,单节点Pandas读取效率低,可以用PySpark分布式处理:
from pyspark.sql import SparkSession from pyspark.sql.functions import input_file_name, lit spark = SparkSession.builder.appName("ParquetStats").getOrCreate() # 递归读取所有Parquet文件,保留文件路径元数据 df = spark.read.parquet("/dbfs/mnt/sgi/report/").withColumn("path", input_file_name()) # 按路径分组统计行数 stats_df = df.groupBy("path").count().withColumnRenamed("count", "count_of_rows") # 统计列数(同个Parquet文件列数一致,取第一行列数即可) col_count = len(df.columns) - 1 # 减去新增的path列 stats_df = stats_df.withColumn("count_of_col", lit(col_count)) # 转Pandas后生成Excel stats_df.toPandas().to_excel("/dbfs/mnt/sgi/report/parquet_stats_spark.xlsx", index=False)
关键说明
- 修正后的Pandas代码通过
paths列表收集每个文件的完整路径,保证路径和行列数一一对应。 - Databricks中生成的Excel文件保存到DBFS路径后,可通过平台界面下载到本地。
- 大文件场景优先用PySpark,避免单节点内存不足的问题。
内容的提问来源于stack exchange,提问作者Anubhav Swaroop
相关产品推荐
相关产品推荐

