Spark中如何提取文件路径中的值作为列添加到DataFrame中
解决方案
Azure Synapse 环境下不管是用 PySpark 还是 Serverless SQL 池,都可以通过内置函数直接获取源文件路径,再按规则提取 State 值即可:
PySpark 场景(Notebook 常用)
用 Spark 内置的 input_file_name() 函数获取当前行对应的源文件全路径,再拆分/正则提取 State 即可,示例代码如下:
from pyspark.sql.functions import input_file_name, split, regexp_extract # 第一步:按你原有的通配符路径读取文件,同时新增列存储源文件全路径 df = spark.read.parquet("abfss://你的容器@你的存储.dfs.core.chinacloudapi.cn/父路径/*/*.parquet") \ .withColumn("source_file_path", input_file_name())
提取 State 两种方式二选一:
- 如果你的路径结构固定,State 位于固定的路径分段,用字符串拆分更简单
举个例子:你的路径格式为/父路径/California/2023/数据文件.parquet,State 位于按/拆分后的第 6 位(索引从0开始计数,你可以先打印几行source_file_path确认对应索引)df_with_state = df.withColumn("State", split("source_file_path", "/")[5]) - 如果路径结构有轻微波动,用正则匹配更稳妥,你可以根据自己的路径规则调整正则表达式
# 示例正则匹配路径中「/ + State名称 + / + 4位年份」的结构,提取第一个分组的内容作为State state_regex = r'\/([A-Za-z\s]+)\/\d{4}\/' df_with_state = df.withColumn("State", regexp_extract("source_file_path", state_regex, 1))
收尾操作
不需要保留源路径的话直接删除临时列即可:
df_final = df_with_state.drop("source_file_path")
Serverless SQL 池场景
如果用 SQL 语句读取文件,用内置的 filepath() 函数即可直接提取通配符对应位置的内容:
SELECT *, -- 括号内的数字对应BULK路径中通配符的序号,比如下面路径中第二个*对应State的位置,就填2 filepath(2) AS State FROM OPENROWSET( BULK 'abfss://你的容器@你的存储.dfs.core.chinacloudapi.cn/父路径/*/*/*.parquet', FORMAT = 'PARQUET' ) AS [result]
内容的提问来源于stack exchange,提问作者Joel Cochran
相关产品推荐
相关产品推荐

