Databricks中Pandas to_csv追加文件失败(OSError Errno95)
问题
尝试通过以下代码实现文件追加功能:创建b.csv文件,每次迭代向文件追加数据,虽然指定了mode='a'参数,但仅能创建文件无法完成追加,执行时抛出OSError: [Errno 95] Operation not supported错误。
代码示例
files = dbutils.fs.ls("/mnt/lake/RAW/test/billion-row-ingestion-time/table/") parquet_file_list = [each.path for each in files if each.name!='_delta_log/'] for each in parquet_file_list: i=0 df = spark.read.parquet(each).toPandas() df.to_csv('/dbfs/FileStore/raw/billion-row-ingestion-time/b.csv', mode='a') print("interation: ", i+1)
报错信息
OSError: [Errno 95] Operation not supported --------------------------------------------------------------------------- OSError Traceback (most recent call last) File /databricks/python/lib/python3.10/site-packages/pandas/io/formats/csvs.py:261, in CSVFormatter.save(self) 251 self.writer = csvlib.writer( 252 handles.handle, 253 lineterminator=self.line_terminator, (...) 258 quotechar=self.quotechar, 259 ) --> 261 self._save() File /databricks/python/lib/python3.10/site-packages/pandas/io/formats/csvs.py:266, in CSVFormatter._save(self) 265 self._save_header() --> 266 self._save_body() File /databricks/python/lib/python3.10/site-packages/pandas/io/formats/csvs.py:304, in CSVFormatter._save_body(self) 303 break --> 304 self._save_chunk(start_i, end_i) File /databricks/python/lib/python3.10/site-packages/pandas/io/formats/csvs.py:315, in CSVFormatter._save_chunk(self, start_i, end_i) 314 ix = self.data_index[slicer]._format_native_types(**self._number_format) --> 315 libwriters.write_csv_rows( 316 data, 317 ix, 318 self.nlevels, 319 self.cols, 320 self.writer, 321 ) File /databricks/python/lib/python3.10/site-packages/pandas/_libs/writers.pyx:55, in pandas._libs.writers.write_csv_rows() OSError: [Errno 95] Operation not supported During handling of the above exception, another exception occurred: OSError Traceback (most recent call last) File <command-1964363491723333>, line 4 2 i=0 3 df = spark.read.parquet(each).toPandas() ----> 4 df.to_csv('/dbfs/FileStore/raw/billion-row-ingestion-time/b.csv', mode='a') 5 print("interation: ", i+1) File /databricks/python/lib/python3.10/site-packages/pandas/core/generic.py:3551, in NDFrame.to_csv(self, path_or_buf, sep, na_rep, float_format, columns, header, index, index_label, mode, encoding, compression, quoting, quotechar, line_terminator, chunksize, date_format, doublequote, escapechar, decimal, errors, storage_options) 3540 df = self if isinstance(self, ABCDataFrame) else self.to_frame() 3542 formatter = DataFrameFormatter( 3543 frame=df, 3544 header=header, (...) 3548 decimal=decimal, 3549 ) -> 3551 return DataFrameRenderer(formatter).to_csv( 3552 path_or_buf, 3553 line_terminator=line_terminator, 3554 sep=sep, 3555 encoding=encoding, 3556 errors=errors, 3557 compression=compression, 3558 quoting=quoting, 3559 columns=columns, 3560 index_label=index_label, 3561 mode=mode, 3562 chunksize=chunksize, 3563 quotechar=quotechar, 3564 date_format=date_format, 3565 doublequote=doublequote, 3566 escapechar=escapechar, 3567 storage_options=storage_options, 3568 ) File /databricks/python/lib/python3.10/site-packages/pandas/io/formats/format.py:1180, in DataFrameRenderer.to_csv(self, path_or_buf, encoding, sep, columns, index_label, mode, compression, quoting, quotechar, line_terminator, chunksize, date_format, doublequote, escapechar, errors, storage_options) 1159 created_buffer = False 1161 csv_formatter = CSVFormatter( 1162 path_or_buf=path_or_buf, 1163 line_terminator=line_terminator, (...) 1178 formatter=self.fmt, 1179 ) -> 1180 csv_formatter.save() 1182 if created_buffer: 1183 assert isinstance(path_or_buf, StringIO) File /databricks/python/lib/python3.10/site-packages/pandas/io/formats/csvs.py:241, in CSVFormatter.save(self) 237 """ 238 Create the writer & save. 239 """ 240 # apply compression and byte/text conversion -> 241 with get_handle( 242 self.filepath_or_buffer, 243 self.mode, 244 encoding=self.encoding, 245 errors=self.errors, 246 compression=self.compression, 247 storage_options=self.storage_options, 248 ) as handles: 249 250 # Note: self.encoding is irrelevant here 251 self.writer = csvlib.writer( 252 handles.handle, 253 lineterminator=self.line_terminator, (...) 258 quotechar=self.quotechar, 259 ) 261 self._save() File /databricks/python/lib/python3.10/site-packages/pandas/io/common.py:124, in IOHandles.__exit__(self, *args) 123 def __exit__(self, *args: Any) -> None: -> 124 self.close() File /databricks/python/lib/python3.10/site-packages/pandas/io/common.py:116, in IOHandles.close(self) 114 self.created_handles.remove(self.handle) 115 for handle in self.created_handles: -> 116 handle.close() 117 self.created_handles = [] 118 self.is_wrapped = False OSError: [Errno 95] Operation not supported
解决方案
问题原因
DBFS(Databricks文件系统)的FUSE挂载路径(/dbfs/...)不支持文件追加操作,这是底层存储协议的限制,因此使用Pandas的to_csv指定mode='a'会抛出不支持操作的错误。
可行方案
方案1:使用Spark API直接合并写入(推荐)
避免循环转换Pandas DataFrame,直接用Spark读取所有Parquet文件合并后一次性写入CSV,效率更高且符合DBFS操作规范:
files = dbutils.fs.ls("/mnt/lake/RAW/test/billion-row-ingestion-time/table/") parquet_file_list = [each.path for each in files if each.name != '_delta_log/'] # 读取所有Parquet文件并合并 df = spark.read.parquet(*parquet_file_list) # 写入CSV,mode='overwrite'覆盖已有文件;若需保留历史数据,可先读取原有CSV合并后再写入 df.write.mode('overwrite').option('header', 'true').csv('/dbfs/FileStore/raw/billion-row-ingestion-time/b.csv')
方案2:本地临时文件中转后上传
利用本地磁盘支持追加的特性,先将数据追加到本地临时文件,完成后再上传到DBFS:
import os local_temp_path = '/tmp/b.csv' target_dbfs_path = '/dbfs/FileStore/raw/billion-row-ingestion-time/b.csv' # 初始化本地文件并写入表头(仅第一次执行) if not os.path.exists(local_temp_path): first_df = spark.read.parquet(parquet_file_list[0]).toPandas() first_df.to_csv(local_temp_path, mode='w', index=False) # 处理剩余文件,不重复写表头 for each in parquet_file_list[1:]: df = spark.read.parquet(each).toPandas() df.to_csv(local_temp_path, mode='a', header=False, index=False) else: # 直接追加所有文件,不写表头 for each in parquet_file_list: df = spark.read.parquet(each).toPandas() df.to_csv(local_temp_path, mode='a', header=False, index=False) # 将本地文件移动到DBFS,覆盖原有文件 dbutils.fs.cp(f'file:{local_temp_path}', target_dbfs_path, overwrite=True) # 清理本地临时文件 os.remove(local_temp_path)
注意事项
- 方案1为最优选择,Spark原生API适配大数据场景,避免Pandas的内存限制和循环写入的低效问题。
- 若必须使用循环追加逻辑,方案2通过本地中转规避DBFS限制,但需确保本地磁盘空间足够容纳全部数据。
内容的提问来源于stack exchange,提问作者crepantherx
相关产品推荐
相关产品推荐

