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

如何将Azure ADLS中指定时间后修改的文件加载到PySpark DataFrame?

按last_modified过滤文件并加载到PySpark

核心问题解答

  1. 获取文件修改时间的函数:Spark 3.0及以上版本提供了input_file_modification_time()函数,和input_file_name()用法一致,会返回文件的毫秒级修改时间戳,可直接用于加载后的DataFrame添加列。

  2. 最优加载方式:优先采用先过滤文件路径再加载数据的方案,避免读取不必要的文件,大幅节省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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 03:33:21