如何使用PySpark按文件名日期过滤加载Databricks数据湖CSV文件
PySpark的load方法本身不支持直接在路径参数中写大于等于的条件筛选文件名,你可以通过以下两种常用方式实现需求:
方案1:预筛选文件路径后加载(推荐,性能最优)
这种方式会先列出目录下的所有文件,筛选出符合日期要求的路径再传给加载方法,避免读取无用文件,适合大数据量场景,Databricks环境可以直接用dbutils.fs.ls实现:
from datetime import datetime # 定义指定的最小日期前缀,示例为2021年10月 min_date_prefix = "202110" # 目标文件存储目录 file_dir = "/mnt/path/" # 筛选符合要求的csv文件路径 valid_files = [] for file in dbutils.fs.ls(file_dir): file_name = file.name # 仅处理命名格式为8位日期+.csv的文件 if file_name.endswith(".csv") and len(file_name) == 12: # 提取前6位年月前缀,如需按天筛选可改为file_name[:8] file_date_prefix = file_name[:6] # 日期格式为连续数字,字符串字典序和时间序一致可直接比较 if file_date_prefix >= min_date_prefix: valid_files.append(file.path) # 仅加载筛选后的文件 df_example = (sqlContext .read .format("com.databricks.spark.csv") .option("delimiter", ";") .option("header","true") .option("inferSchema","true") .option("encoding","windows-1252") .load(valid_files))
方案2:加载后通过文件名过滤
适合数据量较小的场景,不需要提前遍历目录,通过PySpark内置函数提取文件名后过滤:
from pyspark.sql.functions import input_file_name, substring, col min_date_prefix = "202110" df_example = (sqlContext .read .format("com.databricks.spark.csv") .option("delimiter", ";") .option("header","true") .option("inferSchema","true") .option("encoding","windows-1252") .load("/mnt/path/*.csv") # 提取文件名中的6位年月前缀 .withColumn("file_date_prefix", substring(input_file_name(), -10, 6)) # 筛选符合日期要求的文件数据 .filter(col("file_date_prefix") >= min_date_prefix) # 可选:删除辅助计算的列 .drop("file_date_prefix"))
内容的提问来源于stack exchange,提问作者Gonza
相关产品推荐
相关产品推荐

