使用PyArrow引擎将Pandas DataFrame写入S3突破1024分区限制的方法
问题描述
使用PyArrow引擎将Pandas DataFrame按a_id列分区写入S3 Parquet文件时,触发分区数量超过1024的限制错误;尝试直接添加max_partitions参数后,又收到未知参数的报错。
原代码
df.to_parquet(s3_output_path, compression='snappy', engine = 'pyarrow', basename_template = 'part-{i}' + '.parquet', partition_cols = ['a_id'], existing_data_behavior = 'overwrite_or_ignore')
初始错误信息
timestamp,message Traceback (most recent call last): 1727275500645," File ""/opt/predict.py"", line 149, in <module>" 1727275500645, llm.demo() 1727275500645," File ""/opt/predict.py"", line 99, in demo" 1727275500645," df1.to_parquet(s3_output_path, " 1727275500645," File ""/usr/local/lib/python3.10/dist-packages/pandas/util/_decorators.py"", line 207, in wrapper" 1727275500645," return func(*args, **kwargs)" 1727275500645," File ""/usr/local/lib/python3.10/dist-packages/pandas/core/frame.py"", line 2835, in to_parquet" 1727275500646, return to_parquet( 1727275500646," File ""/usr/local/lib/python3.10/dist-packages/pandas/io/parquet.py"", line 420, in to_parquet" 1727275500646, impl.write( 1727275500646," File ""/usr/local/lib/python3.10/dist-packages/pandas/io/parquet.py"", line 186, in write" 1727275500646, self.api.parquet.write_to_dataset( 1727275500646," File ""/usr/local/lib/python3.10/dist-packages/pyarrow/parquet/__init__.py"", line 3153, in write_to_dataset" 1727275500646, ds.write_dataset( 1727275500646," File ""/usr/local/lib/python3.10/dist-packages/pyarrow/dataset.py"", line 930, in write_dataset" 1727275500646, _filesystemdataset_write( 1727275500646," File ""pyarrow/_dataset.pyx"", line 2737, in pyarrow._dataset._filesystemdataset_write"," File ""pyarrow/error.pxi"", line 100, in pyarrow.lib.check_status",pyarrow.lib.ArrowInvalid: Fragment would be written into 34865 partitions. This exceeds the maximum of 1024
添加max_partitions后的错误
File "/usr/local/lib/python3.10/dist-packages/pandas/io/parquet.py", line 186, in write 2024-09-25T15:42:21.674Z self.api.parquet.write_to_dataset( 2024-09-25T15:42:21.674Z File "/usr/local/lib/python3.10/dist-packages/pyarrow/parquet/__init__.py", line 3137, in write_to_dataset 2024-09-25T15:42:21.675Z write_options = parquet_format.make_write_options(**kwargs) 2024-09-25T15:42:21.675Z File "pyarrow/_dataset_parquet.pyx", line 181, in pyarrow._dataset_parquet.ParquetFileFormat.make_write_options 2024-09-25T15:42:21.675Z File "pyarrow/_dataset_parquet.pyx", line 514, in pyarrow._dataset_parquet.ParquetFileWriteOptions.update 2024-09-25T15:42:21.675Z TypeError: unexpected parquet write option: max_partitions
解决方案
max_partitions不是PyArrow Parquet写入选项的参数,而是pyarrow.dataset.write_dataset的专属参数。要在Pandas中传递该参数,需通过engine_kwargs将其转发到底层的write_dataset方法:
修改后的Pandas代码
df.to_parquet( s3_output_path, compression='snappy', engine='pyarrow', basename_template='part-{i}.parquet', partition_cols=['a_id'], existing_data_behavior='overwrite_or_ignore', engine_kwargs={ 'write_dataset_kwargs': { 'max_partitions': 40000 # 设置为大于实际分区数的值,比如这里的34865 } } )
替代方案:直接使用PyArrow API
如果上述方法因版本兼容性问题失效,可以绕过Pandas,直接用PyArrow的API写入数据集:
import pyarrow as pa import pyarrow.parquet as pq from pyarrow import dataset as ds # 将Pandas DataFrame转换为PyArrow Table table = pa.Table.from_pandas(df) # 初始化S3文件系统 fs = pa.fs.S3FileSystem() # 写入数据集并设置分区上限 pq.write_to_dataset( table, root_path=s3_output_path, filesystem=fs, compression='snappy', basename_template='part-{i}.parquet', partition_cols=['a_id'], existing_data_behavior='overwrite_or_ignore', write_dataset_kwargs={'max_partitions': 40000} )
内容的提问来源于stack exchange,提问作者slysid
相关产品推荐
相关产品推荐

