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

Databricks多节点集群无法生成含表头的空DataFrame CSV文件

问题:空DataFrame写入本地磁盘在多节点Databricks集群不生成CSV文件?

编辑说明:无论DataFrame是否为空,此问题均会出现。

通过df.repartition(1).write.save()将空DataFrame写入驱动本地磁盘,设置header=True,期望输出仅含表头行的CSV文件。在Databricks单节点作业集群运行时可生成目标文件,但在2个Worker+1个Driver的多节点作业集群运行相同代码时,仅能看到写入痕迹(如_SUCCESS文件),却未生成CSV文件。

相关代码

import shutil, tempfile, os
from pathlib import Path
from pyspark.sql import SparkSession, DataFrame

def write_as_one_file(spark_session: SparkSession, df: DataFrame, out_dir: str,
                             file_format: str = 'csv', compression: str = 'gzip') -> str:
    """
    @return: full path to the output file
    """
    # When running on Databricks cluster, we need to prefix path with "file:" for it to work.
    prefix = '' if (not spark_session or spark_session.conf.get('spark.app.name') != 'Databricks Shell') else 'file:'
    save_options = {'format': file_format, 'header': True}
    if compression:
        save_options['compression'] = compression

    print(f'prefix: {prefix}, save_options: {save_options}')
    with tempfile.TemporaryDirectory() as tmp_dir_name:
        tmp_df_dir = f'{prefix}{tmp_dir_name}/df'
        print(f'tmp_df_dir: {tmp_df_dir}')
        df.repartition(1).write.save(path=tmp_df_dir, **save_options)
        files = list(Path(tmp_dir_name).rglob(f'part*.{file_format}*'))
        print(f'listing -- tmp_df_dir: {list(Path(tmp_dir_name).rglob("*"))}')
        if len(files) != 1:
            raise Exception(f'expected exactly only one file to match pattern, found "{files}"')

        os.makedirs(out_dir, exist_ok=True)
        out_file = f'{out_dir}/{os.path.basename(files[0])}'
        shutil.copy(files[0], out_file)
        return out_file

df1 = spark.createDataFrame(data=[], schema='id:int')
write_as_one_file(spark, df1, '/home/kash/output_dir/')

运行日志对比

单节点集群(正常生成文件)

prefix: file:, save_options: {'format': 'csv', 'header': True, 'compression': 'gzip'}
tmp_df_dir: file:/tmp/tmp3xl8rvi7/df
listing -- tmp_df_dir: [
    PosixPath('/tmp/tmp3xl8rvi7/df'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/_SUCCESS'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/.part-00000-tid-1201247493872524145-a40b5e0b-2dde-4edc-98ec-c40cbbf6a29d-1-1-c000.csv.gz.crc'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/part-00000-tid-1201247493872524145-a40b5e0b-2dde-4edc-98ec-c40cbbf6a29d-1-1-c000.csv.gz'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/_committed_1201247493872524145'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/._SUCCESS.crc'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/._committed_1201247493872524145.crc'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/_started_1201247493872524145'), 
    PosixPath('/tmp/tmp3xl8rvi7/df/._started_1201247493872524145.crc')
]

多节点集群(未生成CSV文件)

prefix: file:, save_options: {'format': 'csv', 'header': True, 'compression': 'gzip'}
tmp_df_dir: file:/tmp/tmpj2ki4hpl/df
listing -- tmp_df_dir: [
    PosixPath('/tmp/tmpj2ki4hpl/df'), 
    PosixPath('/tmp/tmpj2ki4hpl/df/._committed_637987277244855410.crc'), 
    PosixPath('/tmp/tmpj2ki4hpl/df/_SUCCESS'), 
    PosixPath('/tmp/tmpj2ki4hpl/df/_committed_637987277244855410'), 
    PosixPath('/tmp/tmpj2ki4hpl/df/._SUCCESS.crc')
]

# 此处抛出异常...

环境信息

  • Databricks Runtime 13.3LTS(Spark 3.4.1),开启/关闭Photon无影响
  • 单节点集群节点类型:r6id.xlarge
  • 多节点集群节点类型:i3.xlarge

原因分析

  1. 集群节点文件系统隔离:多节点集群中,repartition(1)后的分区任务会被Spark调度到任意Worker节点执行。你写入的路径是Driver本地的临时目录(/tmp/...),而Worker节点的本地文件系统与Driver完全隔离,生成的part*.csv.gz文件会被写到Worker的本地磁盘,Driver端自然找不到该文件。
  2. 单节点集群无隔离问题:单节点集群中Worker与Driver是同一个节点,任务执行后文件直接写入Driver本地目录,因此能正常读取。

解决办法

方法一:强制任务在Driver节点执行

通过coalesce(1)结合本地 checkpoint,强制Spark将分区任务调度到Driver执行,确保文件写入Driver本地磁盘:

def write_as_one_file(spark_session: SparkSession, df: DataFrame, out_dir: str,
                             file_format: str = 'csv', compression: str = 'gzip') -> str:
    # ... 原有前缀、保存配置逻辑不变 ...
    with tempfile.TemporaryDirectory() as tmp_dir_name:
        tmp_df_dir = f'{prefix}{tmp_dir_name}/df'
        # 强制任务在Driver执行
        df_local = df.coalesce(1).localCheckpoint()
        df_local.write.save(path=tmp_df_dir, **save_options)
        # ... 后续文件查找、复制逻辑不变 ...

方法二:使用DBFS作为中转存储

先将文件写入Databricks分布式文件系统(DBFS),再从DBFS的本地挂载点(/dbfs/)复制到Driver本地目录,避免节点文件系统隔离问题:

def write_as_one_file(spark_session: SparkSession, df: DataFrame, out_dir: str,
                             file_format: str = 'csv', compression: str = 'gzip') -> str:
    # ... 原有前缀、保存配置逻辑不变 ...
    with tempfile.TemporaryDirectory() as tmp_dir_name:
        # 写入DBFS路径
        tmp_df_dir = f'dbfs:/tmp/{tmp_dir_name}/df'
        df.repartition(1).write.save(path=tmp_df_dir, **save_options)
        # 从DBFS本地挂载点读取文件
        dbfs_local_path = f'/dbfs/tmp/{tmp_dir_name}/df'
        files = list(Path(dbfs_local_path).rglob(f'part*.{file_format}*'))
        # ... 后续文件复制逻辑不变 ...

方法三:手动生成表头文件(空DataFrame场景)

如果仅处理空DataFrame的情况,可直接手动生成含表头的文件,完全绕开Spark的分区调度逻辑:

import gzip

def write_as_one_file(spark_session: SparkSession, df: DataFrame, out_dir: str,
                             file_format: str = 'csv', compression: str = 'gzip') -> str:
    os.makedirs(out_dir, exist_ok=True)
    # 空DataFrame直接生成表头文件
    if df.isEmpty():
        header = ','.join(df.columns)
        out_file = f'{out_dir}/header_only.csv.gz'
        with gzip.open(out_file, 'wt') as f:
            f.write(header + '\n')
        return out_file
    # ... 原有非空DataFrame处理逻辑 ...

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 22:50:55