PySpark读取Parquet时如何跳过不存在路径避免报错
HDFS 路径 /dir1/dir2/ 下存储了多份 Parquet 文件,文件名包含随机生成的时间戳,示例文件路径为 /dir1/dir2/2022-06-16-03-12-36-086.snappy.parquet,文件名各字段含义如下:
2022、06、16分别对应年、月、日03、12、36、086分别对应时、分、秒、毫秒
需求为读取时间戳介于 2022-06-16-04-15-00-000 到 2022-06-16-05-15-00-000 之间的所有文件,初始实现代码如下:
paths = [f'/dir1/dir2/{tm.date()}-{tm.hour:02d}-{tm.minute:02d}-*' \ for tm in pd.date_range('2022-06-16 04:15:00','2022-06-16 05:15:00', freq = 'min')] df = spark.read.parquet(*paths)
由于不是所有分钟、秒、毫秒维度都存在对应文件,执行代码时抛出如下错误:
AnalysisException: Path does not exist: /dir1/dir2/2022-06-16-04-23-*.snappy.parquet
尝试添加 Spark 配置 ("spark.sql.files.ignoreMissingFiles", "true") 后错误仍然存在,需要实现仅读取实际存在的路径,自动跳过不存在的路径不触发报错。
spark.sql.files.ignoreMissingFiles 参数仅在作业运行阶段生效,处理的是路径解析完成后、任务读取数据时文件被意外删除的场景,不会在路径解析阶段跳过没有任何匹配结果的通配路径,因此无法解决当前报错。
方案1:读取前过滤有效路径
调用 Hadoop FileSystem API 提前校验生成的通配路径是否有匹配文件,只将真实存在的路径传入 Spark 读取方法,代码示例:
from py4j.java_gateway import java_import # 导入Hadoop Path类 java_import(spark._jvm, 'org.apache.hadoop.fs.Path') # 初始化HDFS客户端 fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) # 生成原始通配路径列表 raw_paths = [f'/dir1/dir2/{tm.date()}-{tm.hour:02d}-{tm.minute:02d}-*' \ for tm in pd.date_range('2022-06-16 04:15:00','2022-06-16 05:15:00', freq = 'min')] # 过滤出实际存在匹配文件的路径 valid_paths = [] for path in raw_paths: if fs.globStatus(spark._jvm.Path(path)): valid_paths.append(path) # 传入有效路径读取数据 df = spark.read.parquet(*valid_paths)
该方案只扫描目标时间范围内的路径元数据,适合目录下总文件量较大的场景。
方案2:读取父目录后按文件名过滤
不按分钟粒度拼接通配路径,直接读取上层父目录,通过内置函数提取文件名中的时间戳做范围过滤,代码示例:
from pyspark.sql import functions as F df = spark.read.parquet("/dir1/dir2/") \ # 从文件全路径中提取时间戳字符串 .withColumn("file_timestamp", F.regexp_extract( F.input_file_name(), r'(\d{4}-\d{2}-\d{2}-\d{2}-\d{2}-\d{2}-\d{3})\.snappy\.parquet$', 1 )) \ # 过滤指定时间范围的文件 .filter(F.col("file_timestamp").between( "2022-06-16-04-15-00-000", "2022-06-16-05-15-00-000" ))
该方案写法更简洁,但会先扫描整个目录下所有文件的元数据,适合目录总文件量适中的场景。
内容的提问来源于stack exchange,提问作者aishik roy chaudhury

