如何将Parquet文件的创建日期添加至Spark DataFrame
为Spark DataFrame添加Parquet文件的创建日期列
针对你的需求,这里提供几种适配Databricks环境的解决方案,避免使用本地os模块(分布式场景下无法跨节点访问文件元数据):
方案一:通过文件路径关联元数据(高效,适合大数据量)
利用Spark的input_file_name()获取每行数据对应的源文件路径,再提前收集所有文件的创建时间并关联到DataFrame:
from pyspark.sql.functions import input_file_name, lit from datetime import datetime # 1. 遍历所有日期文件夹,收集Parquet文件路径和对应的创建时间 file_info = [] date_folders = dbutils.fs.ls("/mnt/dev/bronze/Voucher/") for folder in date_folders: # 遍历当前日期文件夹下的文件 parquet_files = dbutils.fs.ls(folder.path) for file in parquet_files: if file.path.endswith(".parquet"): # 将毫秒级时间戳转换为datetime格式 create_time = datetime.fromtimestamp(file.modificationTime / 1000) file_info.append( (file.path, create_time) ) # 2. 将文件信息转换为Spark DataFrame,用于后续关联 file_meta_df = spark.createDataFrame(file_info, ["file_path", "create_date"]) # 3. 读取原始数据并添加源文件路径列 raw_df = spark.read.parquet("/mnt/dev/bronze/Voucher/*/*").withColumn("file_path", input_file_name()) # 4. 关联元数据,添加创建日期列并移除路径列 final_df = raw_df.join(file_meta_df, on="file_path", how="left").drop("file_path")
方案二:遍历单个文件读取并合并(适合小数据量)
逐个读取每个Parquet文件,直接添加创建日期列后合并所有DataFrame:
from pyspark.sql import DataFrame from datetime import datetime all_dfs = [] date_folders = dbutils.fs.ls("/mnt/dev/bronze/Voucher/") for folder in date_folders: parquet_files = dbutils.fs.ls(folder.path) for file in parquet_files: if file.path.endswith(".parquet"): create_time = datetime.fromtimestamp(file.modificationTime / 1000) # 读取单个文件并添加创建日期列 single_df = spark.read.parquet(file.path).withColumn("create_date", lit(create_time)) all_dfs.append(single_df) # 合并所有DataFrame(保证列名一致) final_df = all_dfs[0] for df in all_dfs[1:]: final_df = final_df.unionByName(df)
方案三:从文件夹路径提取日期(最简,需路径日期与创建日期一致)
如果你的日期子文件夹命名(如2022-09-23)和文件实际创建日期完全一致,可以直接从文件路径提取日期,无需读取文件元数据:
from pyspark.sql.functions import regexp_extract, input_file_name, to_date raw_df = spark.read.parquet("/mnt/dev/bronze/Voucher/*/*") final_df = raw_df.withColumn( "create_date", # 正则提取路径中的日期部分 regexp_extract(input_file_name(), r"/Voucher/(\d{4}-\d{2}-\d{2})/", 1) ).withColumn("create_date", to_date("create_date"))
注意事项
- 避免使用
os.stat:在Databricks分布式环境中,os模块只能访问Driver节点的本地文件,无法正确获取DBFS上的文件元数据,推荐使用dbutils.fs.ls来遍历文件并获取元信息。 - 性能选择:大数据量场景优先使用方案一,小数据量可以用方案二,路径日期可靠时优先用方案三。
内容的提问来源于stack exchange,提问作者Salah K.
相关产品推荐
相关产品推荐

