PyArrow fsspec缓存S3 Parquet分区后修改引发FileNotFoundError问题
PyArrow + fsspec读取S3 Parquet分区时缓存导致新增/修改文件触发FileNotFoundError的解决方案
问题现象
使用pandas.read_parquet(底层基于PyArrow Dataset)读取S3上的Parquet分区后,若修改该分区内容(新增文件或覆盖已有文件),再次读取时会触发异常:
- 新增文件时,读取提示找不到新上传的文件
- 覆盖已有文件时,提示原文件ETag失效、文件不存在
复现代码
import os import boto3 import pandas as pd df1 = pd.DataFrame([{'a': i, 'b': i} for i in range(10)]) df1.to_parquet('part1.parquet') df2 = pd.DataFrame([{'a': i+100, 'b': i+100} for i in range(10)]) df2.to_parquet('part2.parquet') s3_client = boto3.client('s3') url = 's3://bucket_name/test01' s3_client.upload_file('part1.parquet', 'bucket_name', os.path.join('test01', 'p=x', 'part1.parquet')) dfx1 = pd.read_parquet(url) # 第一次读取正常 # 新增文件场景 s3_client.upload_file('part2.parquet', 'bucket_name', os.path.join('test01', 'p=x', 'part2.parquet')) dfx2 = pd.read_parquet(url) # 触发FileNotFoundError # 覆盖文件场景(替换上述新增步骤) # s3_client.upload_file('part2.parquet', 'bucket_name', os.path.join('test01', 'p=x', 'part1.parquet')) # dfx2 = pd.read_parquet(url) # 触发FileExpired异常
新增文件场景异常回溯
File "<env>/lib/python3.8/site-packages/pandas/io/parquet.py", line 493, in read_parquet return impl.read( File "<env>/lib/python3.8/site-packages/pandas/io/parquet.py", line 240, in read result = self.api.parquet.read_table( File "<env>/lib/python3.8/site-packages/pyarrow/parquet.py", line 1996, in read_table return dataset.read(columns=columns, use_threads=use_threads, File "<env>/lib/python3.8/site-packages/pyarrow/parquet.py", line 1831, in read table = self._dataset.to_table( File "pyarrow/_dataset.pyx", line 323, in pyarrow._dataset.Dataset.to_table File "pyarrow/_dataset.pyx", line 2311, in pyarrow._dataset.Scanner.to_table File "pyarrow/error.pxi", line 143, in pyarrow.lib.pyarrow_internal_check_status File "pyarrow/_fs.pyx", line 1179, in pyarrow._fs._cb_open_input_file File "<env>/lib/python3.8/site-packages/pyarrow/fs.py", line 394, in open_input_file raise FileNotFoundError(path) FileNotFoundError: bucket_name/test01/p=x/part2.parquet
覆盖文件场景异常回溯
File "pyarrow/_dataset.pyx", line 1680, in pyarrow._dataset.DatasetFactory.finish File "pyarrow/error.pxi", line 143, in pyarrow.lib.pyarrow_internal_check_status File "<env>/lib/python3.8/site-packages/fsspec/spec.py", line 1544, in read out = self.cache._fetch(self.loc, self.loc + length) File "/<env>/lib/python3.8/site-packages/fsspec/caching.py", line 377, in _fetch self.cache = self.fetcher(start, bend) File "<env>/lib/python3.8/site-packages/s3fs/core.py", line 1965, in _fetch_range raise FileExpired( s3fs.utils.FileExpired: [Errno 16] The remote file corresponding to filename bucket_name/test01/p=x/part1.parquet and Etag "d64a2de4f9c93dff49ecd3f19c414f61" no longer exists.
注:在新Python进程中读取该分区无异常,问题源于当前进程内的缓存未更新。
原因分析
- PyArrow Dataset缓存:第一次读取时,PyArrow会缓存S3分区的文件列表,后续读取直接使用缓存列表,不会重新扫描S3获取最新文件信息,导致新增文件无法被识别。
- fsspec缓存:fsspec会缓存S3文件的元数据(如ETag)和内容,文件被覆盖后,缓存元数据与远程文件不匹配,触发
FileExpired异常。
重载pandas或PyArrow模块无法清除这些缓存,因为缓存存储在模块级别的对象中。
解决方案
方案1:禁用PyArrow的文件系统缓存
显式创建带缓存禁用配置的S3文件系统,传给read_parquet的filesystem参数,强制每次读取重新扫描S3文件列表:
import pyarrow.fs # 创建禁用缓存的S3文件系统 fs = pyarrow.fs.S3FileSystem(default_cache_options={'enabled': False}) dfx2 = pd.read_parquet(url, filesystem=fs)
方案2:重置fsspec缓存
调用fsspec的缓存失效方法,清除当前进程内的S3文件元数据缓存:
import fsspec # 重置指定路径的缓存 fsspec.invalidate_cache(url) # 或者重置所有缓存 # fsspec.invalidate_cache() dfx2 = pd.read_parquet(url)
方案3:使用fsspec创建无缓存的文件系统
直接使用fsspec的S3FileSystem并禁用缓存,传给read_parquet:
from fsspec.implementations.s3 import S3FileSystem # 创建无缓存的S3文件系统 fs = S3FileSystem(cache_type='none') dfx2 = pd.read_parquet(url, filesystem=fs)
方案4:针对覆盖文件的额外处理
若存在文件覆盖场景,可在读取时指定refresh=True(fsspec参数),强制刷新文件元数据:
dfx2 = pd.read_parquet(url, storage_options={'refresh': True})
内容的提问来源于stack exchange,提问作者michaelgbj
相关产品推荐
相关产品推荐

