Azure环境下将PySpark DataFrame保存至ADLS动态文件夹并生成单CSV
Azure Databricks 导出PySpark DataFrame到ADLS动态路径并生成单个指定CSV文件
核心需求
- 将PySpark DataFrame写入Azure Data Lake(ADLS)
- 路径需按「当前年/月/日」动态生成(如
abfss://container@storageaccount.dfs.core.windows.net/output/2024/05/20/) - 生成单个指定名称的CSV文件(如
abc.csv)
解决方案步骤
1. 确保ADLS访问配置正确
如果使用DBFS挂载:
执行以下代码完成ADLS挂载(以Service Principal为例),挂载后DBFS路径直接映射到ADLS,写入挂载路径即写入ADLS:
configs = {"fs.azure.account.auth.type": "OAuth", "fs.azure.account.oauth.provider.type": "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider", "fs.azure.account.oauth2.client.id": "<your-client-id>", "fs.azure.account.oauth2.client.secret": "<your-client-secret>", "fs.azure.account.oauth2.client.endpoint": "https://login.microsoftonline.com/<your-tenant-id>/oauth2/token"} # 挂载ADLS容器到DBFS路径 dbutils.fs.mount( source = "abfss://<container-name>@<storage-account-name>.dfs.core.windows.net/", mount_point = "/mnt/adls", extra_configs = configs)
如果直接使用ABFS路径,确保Databricks集群已配置Service Principal权限,无需挂载即可直接访问abfss://container@storageaccount.dfs.core.windows.net/路径。
2. 生成动态日期路径
用Python获取当前年月日,拼接目标路径:
from datetime import datetime current_date = datetime.now() # 按年/月/日格式生成路径(挂载版) dynamic_path = f"/mnt/adls/output/{current_date.year}/{current_date.month:02d}/{current_date.day:02d}/" # 如果用ABFS直接路径,替换为: # dynamic_path = f"abfss://container@storageaccount.dfs.core.windows.net/output/{current_date.year}/{current_date.month:02d}/{current_date.day:02d}/"
3. 写入单个指定名称CSV文件
PySpark默认按分区生成多个文件,需先合并分区,再通过重命名实现指定文件名:
# 假设你的PySpark DataFrame名为df target_file_name = "abc.csv" temp_path = dynamic_path + "temp/" # 合并为1个分区后写入临时路径,指定header和格式 df.coalesce(1) \ .write \ .mode("overwrite") \ .option("header", "true") \ .csv(temp_path) # 获取临时路径下的part文件(PySpark生成的单个分区文件) part_file = [file.path for file in dbutils.fs.ls(temp_path) if file.path.endswith(".csv")][0] # 重命名part文件为目标文件名,并移动到最终路径 dbutils.fs.mv(part_file, dynamic_path + target_file_name) # 删除临时路径 dbutils.fs.rm(temp_path, recurse=True)
关键说明
- 使用
coalesce(1)而非repartition(1):前者直接合并现有分区,避免数据shuffle,性能更优;如果DataFrame本身分区数多,coalesce无法合并时再考虑repartition。 - 临时路径的作用:PySpark无法直接指定输出文件名,只能先写入临时目录,再将生成的单个part文件重命名为目标文件。
- 模式选择:
.mode("overwrite")会覆盖目标路径下的现有文件,可根据需求改为append或ignore。
补充:如果想用Pandas DataFrame实现动态路径
若数据量较小,也可以结合动态路径使用Pandas:
# 将PySpark DataFrame转为Pandas pdf = df.toPandas() # 确保动态路径存在,不存在则创建 dbutils.fs.mkdirs(dynamic_path) # 写入指定文件名(挂载路径需加/dbfs前缀) pdf.to_csv(f"/dbfs{dynamic_path}{target_file_name}", index=False, header=True) # 或者直接用ABFS路径(需确保集群有权限) # pdf.to_csv(f"{dynamic_path}{target_file_name}", index=False, header=True)
内容的提问来源于stack exchange,提问作者Shikha
相关产品推荐
相关产品推荐

