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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 12:54:04