如何使用PySpark查找指定路径下的最新修改文件?
用PySpark自动读取最新修改的CSV文件
当然可以实现!我来给你分享两种实用的方法,适配不同的场景:
方法一:用Hadoop FileSystem API(推荐,适配分布式存储)
这个方法利用PySpark底层的Hadoop文件系统工具,能兼容本地磁盘、HDFS、S3等多种存储环境,是分布式场景下最稳妥的方案:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("FetchLatestCSV").getOrCreate() sc = spark.sparkContext # 替换成你的目标路径 target_path = "Path://to/file" # 获取Hadoop文件系统实例 hadoop_fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(sc._jsc.hadoopConfiguration()) # 列出路径下的所有文件/目录状态 file_status_list = hadoop_fs.listStatus(sc._jvm.org.apache.hadoop.fs.Path(target_path)) # 筛选出文件(排除目录),收集每个文件的修改时间和路径 file_details = [] for status in file_status_list: if not status.isDirectory(): file_full_path = status.getPath().toString() modify_time = status.getModificationTime() # 时间戳格式 file_details.append((modify_time, file_full_path)) # 按修改时间倒序排序,取第一个就是最新文件 latest_csv = sorted(file_details, key=lambda x: x[0], reverse=True)[0][1] # 读取最新CSV df = spark.read.csv(latest_csv, header=True, inferSchema=True)
方法二:结合Python的os模块(仅适用于本地文件系统)
如果你的文件都在Spark集群的本地磁盘,也可以用Python原生的os模块先找到最新文件,再传给Spark读取:
import os from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ReadLocalLatestCSV").getOrCreate() target_dir = "/local/path/to/files" # 获取目录下所有CSV文件的路径 csv_files = [os.path.join(target_dir, f) for f in os.listdir(target_dir) if f.endswith('.csv')] # 按文件修改时间排序,取最新的 latest_csv = max(csv_files, key=os.path.getmtime) # 读取文件 df = spark.read.csv(latest_csv, header=True, inferSchema=True)
注意事项
- 如果你的路径包含通配符(比如
path/to/*.csv),记得在筛选文件时加上后缀判断,避免读取到无关文件 - 分布式环境下优先用方法一,因为
os模块只能访问Driver节点的本地文件,无法读取HDFS/S3等分布式存储的文件 - 如果需要读取最近N个文件,只需要修改排序后的切片逻辑(比如
sorted(...)[:3]取前3个最新的)
内容的提问来源于stack exchange,提问作者Eles
相关产品推荐
相关产品推荐

