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

如何在Databricks中使用PySpark写入单个CSV文件

Azure Databricks 下PySpark输出单个CSV/Parquet文件实现

大家好!
之前我在Azure Databricks环境中找PySpark写入单个CSV文件的现成方案时,没找到能直接用的实现,于是自己写了对应的工具函数,现在把方案分享出来,也欢迎大家交流补充其他实现思路。

说明:原代码注释为西班牙语,函数核心逻辑如下:

  • 调用coalesce(1)将目标DataFrame合并为单分区,写入指定路径下temp_开头的临时目录
  • 提取临时目录中Spark生成的part-前缀的实际数据文件,按自定义名称重命名后移动到目标存储路径
  • 删除临时目录,最终得到单个目标文件

实现代码如下:

def escribe_fichero_unico(dataframe, path, file_name, file_format = 'csv'):
    """
    功能:
    1. 创建临时目录存放Spark默认写出的分区文件
    2. 合并所有分区为单个独立文件
    3. 将合并后的文件移动到目标路径
    4. 清理临时目录
    参数:
        dataframe: 要保存为单文件的DataFrame
        file_name: 字符串类型,输出文件名称
        file_format: 字符串类型,支持'csv'/'parquet',默认值为'csv'
        path: 字符串类型,文件存储的目标路径
    """
    import os
    
    # 写入单分区数据到临时目录
    path_temp = path + 'temp_' + file_name + '_trash'
    if file_format == 'csv': 
        dataframe.coalesce(1).write.format('csv').mode('overwrite')                \
                            .options(header="true", delimiter=";")  \
                            .save(path_temp)
    elif file_format == 'parquet':
        dataframe.coalesce(1).write.format("parquet").mode("overwrite")  \
                            .save(path_temp)
        
    # 定位临时目录中的实际数据文件
    file_part = [file.path for file in dbutils.fs.ls(path_temp) if os.path.basename(file.path).startswith('part')][0]
    # 移动重命名文件到目标路径
    dbutils.fs.mv(file_part, path + file_name + '.' + file_format)
    # 删除临时目录
    dbutils.fs.rm(path_temp, True)

注意:如果数据量非常大,不建议使用coalesce(1),会把所有数据拉到单个Executor节点处理,容易出现OOM问题,该方案更适合中小规模数据集的单文件导出场景。

希望这个方案能帮到有需要的人 :)


内容的提问来源于stack exchange,提问作者iikelfarina

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 21:09:18