Spark优化:用通配符筛选文件并避免加载过期文件
Spark高效加载指定格式文件(排除过期文件)
需求说明
- 目标文件:以
Products开头、后缀带6位日期戳(如Products_202410)的CSV文件 - 排除文件:
Products_expired_开头的所有过期文件 - 核心限制:文件行数超1亿,禁止用
input_file_name()加载后过滤(成本过高);需避免高成本循环/追加操作,直接精准加载有效文件(比如示例中的Products_202409、Products_202410、Products_202411)
原参考代码问题
原代码逻辑不合理,先尝试加载过期文件判断是否存在,不仅触发不必要的IO,还没解决精准加载有效文件的核心问题:
if spark.read.load("abfss://abc-def-g-001@testing123.dfs.core.windows.net/data/testing/PRODUCTS_Expired_*.csv") contains "expired": mssparkutils.notebook.exit(send_email("abcdefg@gmail.com", "expired", "failed")) else: df = spark.read.load("abfss://abc-def-g-001@testing123.dfs.core.windows.net/data/testing/PRODUCTS_*.csv")
优化方案
方案1:通配符精准匹配(最优性能)
利用Spark支持的正则通配符,直接在加载阶段过滤掉过期文件,无需额外处理:
# 匹配以Products_开头、后跟6位数字日期戳的CSV文件,自动排除Products_expired_开头的文件 target_path = "abfss://abc-def-g-001@testing123.dfs.core.windows.net/data/testing/Products_[0-9]{6}.csv" df = spark.read.load(target_path)
方案2:先获取文件列表再过滤(适配复杂规则)
如果通配符无法满足需求,用mssparkutils.fs.ls先遍历路径获取文件列表,过滤出有效文件后再加载,同样避免无效IO:
# 获取路径下所有文件信息 all_files = mssparkutils.fs.ls("abfss://abc-def-g-001@testing123.dfs.core.windows.net/data/testing/") # 过滤出符合要求的有效文件路径 valid_file_paths = [ file.path for file in all_files if file.name.startswith("Products_") and not file.name.startswith("Products_expired_") ] # 直接加载筛选后的文件 df = spark.read.load(valid_file_paths)
方案优势
- 两种方案都在加载前完成过滤,完全避免加载过期文件,减少大量不必要的IO和计算开销
- 无需加载后调用
input_file_name()处理,规避了超大数据量下的性能损耗 - 没有循环追加操作,直接批量加载,符合Spark分布式处理的最优实践
内容的提问来源于stack exchange,提问作者Dan Wang
相关产品推荐
相关产品推荐

