如何在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
相关产品推荐
相关产品推荐

