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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 00:07:29