如何将Azure ADLS中指定时间后修改的文件加载到PySpark DataFrame?
按last_modified过滤文件并加载到PySpark
核心问题解答
获取文件修改时间的函数:Spark 3.0及以上版本提供了
input_file_modification_time()函数,和input_file_name()用法一致,会返回文件的毫秒级修改时间戳,可直接用于加载后的DataFrame添加列。最优加载方式:优先采用先过滤文件路径再加载数据的方案,避免读取不必要的文件,大幅节省IO和计算资源;小数据集场景下可选择加载后再过滤。
方法一:先过滤文件路径再加载(推荐)
通过Hadoop FileSystem API批量获取文件元数据,筛选出修改时间符合要求的文件后再读取,适合大数据量场景:
from pyspark.sql import SparkSession from datetime import datetime import time # 初始化SparkSession spark = SparkSession.builder.appName("FilterFilesByModifiedTime").getOrCreate() # 定义目标时间(示例:2024年1月1日0点),转换为毫秒级时间戳 target_datetime = datetime(2024, 1, 1, 0, 0, 0) target_timestamp = int(time.mktime(target_datetime.timetuple()) * 1000) # 获取ADLS路径下的所有文件状态 hadoop_conf = spark._jsc.hadoopConfiguration() fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) base_path = spark._jvm.org.apache.hadoop.fs.Path( "abfss://container@storageaccount.dfs.core.windows.net/*/*/*/*/*.json" ) file_status_list = fs.listStatus(base_path) # 筛选出修改时间晚于目标时间的文件路径 filtered_file_paths = [ status.getPath().toString() for status in file_status_list if status.getModificationTime() > target_timestamp ] # 加载筛选后的文件 if filtered_file_paths: df = spark.read.json(filtered_file_paths) # 可选:添加文件名和修改时间列 from pyspark.sql.functions import input_file_name, input_file_modification_time df = df.withColumn("filename", input_file_name()) \ .withColumn("last_modified", input_file_modification_time()) else: print("未找到符合修改时间条件的文件")
方法二:加载后过滤(适合小数据集)
如果文件总数较少,可先加载所有数据,再通过input_file_modification_time()过滤:
from pyspark.sql.functions import input_file_name, input_file_modification_time from datetime import datetime import time # 目标时间转毫秒级时间戳 target_datetime = datetime(2024, 1, 1, 0, 0, 0) target_timestamp = int(time.mktime(target_datetime.timetuple()) * 1000) # 加载所有文件并添加元数据列 df = spark.read.json( "abfss://container@storageaccount.dfs.core.windows.net/*/*/*/*/*.json" ).withColumn("filename", input_file_name()) \ .withColumn("last_modified", input_file_modification_time()) # 过滤出修改时间符合要求的数据 filtered_df = df.filter(df.last_modified > target_timestamp)
注意事项
input_file_modification_time()仅在Spark 3.0+版本可用;若使用旧版Spark,只能通过方法一获取文件修改时间。- 方法一的性能优势在文件数量多、单文件体积大的场景下会非常明显,能有效减少不必要的数据读取。
内容的提问来源于stack exchange,提问作者Aleksei Zhukov
相关产品推荐
相关产品推荐

