ds.read_file()是否全量加载文件至内存?能否支持PyArrow谓词下推?
问题解答
核心结论
ds.read_file()会将整个S3文件下载到本地内存,返回的是类似BytesIO的内存缓冲区,并不支持行组级的懒访问。这意味着你代码里的filters=参数根本没法实现谓词下推——因为所有数据已经被提前加载到内存,过滤操作只是在本地内存中进行,完全无法减少从S3下载的数据量,这就是导致OOM错误的直接原因。
关键依据
你观察到的file_handle.seek(0)操作是核心线索:只有内存中的缓冲区(比如BytesIO)才需要通过seek重置读取位置,而支持流式/懒加载的远程文件句柄(比如PyArrow直接对接S3的句柄)不需要这个操作,因为它们是直接和远程存储交互,不会把全量数据加载到内存。
解决办法
放弃使用ds.read_file(),直接通过PyArrow对接S3存储,结合Dash Enterprise的数据源凭证实现安全访问,这样就能利用PyArrow的原生能力实现谓词下推和懒加载:
from dash_enterprise_libraries import data_sources as ds import pyarrow.parquet as pq import pyarrow.fs # 从Dash Enterprise数据源获取S3凭证 s3_credentials = ds.get_credentials(DATA_SOURCE) # 初始化PyArrow的S3文件系统 s3_filesystem = pyarrow.fs.S3FileSystem( access_key=s3_credentials['access_key_id'], secret_key=s3_credentials['secret_access_key'] # 若你的S3存储指定了区域,添加参数:region=s3_credentials['region'] ) # 直接读取S3上的Parquet文件,自动实现谓词下推 table = pq.read_table( path, columns=columns, filters=filter_expr, filesystem=s3_filesystem )
这种方式下,PyArrow会先读取Parquet文件的元数据,识别出符合过滤条件的行组,仅下载这些行组的数据,大幅降低内存占用,从根源上避免OOM错误。
内容的提问来源于stack exchange,提问作者Insane Gamer
相关产品推荐
相关产品推荐

