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

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进程中读取该分区无异常,问题源于当前进程内的缓存未更新。

原因分析

  1. PyArrow Dataset缓存:第一次读取时,PyArrow会缓存S3分区的文件列表,后续读取直接使用缓存列表,不会重新扫描S3获取最新文件信息,导致新增文件无法被识别。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 14:47:38