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(Spark3.4.1),开启/关闭Photon无影响 - 单节点集群节点类型:
r6id.xlarge - 多节点集群节点类型:
i3.xlarge
原因分析
- 集群节点文件系统隔离:多节点集群中,
repartition(1)后的分区任务会被Spark调度到任意Worker节点执行。你写入的路径是Driver本地的临时目录(/tmp/...),而Worker节点的本地文件系统与Driver完全隔离,生成的part*.csv.gz文件会被写到Worker的本地磁盘,Driver端自然找不到该文件。 - 单节点集群无隔离问题:单节点集群中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
相关产品推荐
相关产品推荐

