【发布时间】:2017-03-03 19:38:26
【问题描述】:
我尝试从 s3 读取大量 csv 文件,工作人员在具有正确 IAM 角色的 ec2 实例上运行(我可以从其他脚本的相同存储桶中读取)。 当我尝试使用此命令从私有存储桶中读取自己的数据时:
client = Client('scheduler-on-ec2')
df = read_csv('s3://xyz/*csv.gz',
compression='gzip',
blocksize=None,
#storage_options={'key': '', 'secret': ''}
)
df.size.compute()
数据看起来像是在本地读取(由本地 Python 解释器,而不是工作人员),然后由本地解释器发送给工作人员(或调度程序?),当工作人员收到块时,他们运行计算并返回结果。通过storage_options 传递或不传递密钥和秘密都相同。
当我使用storage_options={'anon': True} 从公共 s3 存储桶(纽约出租车数据)中读取数据时,一切正常。
您认为问题是什么,我应该重新配置更改以使工作人员直接从 s3 读取?
s3fs 安装正确,根据 dask 支持的文件系统如下:
>>>> dask.bytes.core._filesystems
{'file': dask.bytes.local.LocalFileSystem,
's3': dask.bytes.s3.DaskS3FileSystem}
更新
在监控网络接口后,似乎有些东西从解释器上传到了调度器。数据帧(或包)中的分区越多,发送到调度程序的数据就越大。我以为它可能是计算图,但它确实很大。对于 12 个文件,它是 2-3MB,对于 30 个文件,它是 20MB,对于更大的数据,(150 个文件)将它发送到调度程序需要太长时间,我没有等待它。还有什么被发送到可以占用这么多数据的调度程序?
【问题讨论】:
-
> 还有什么发送到调度程序可以占用这么多数据?我所知道的。没有什么。如果您可以生成可重现的minimal failing example,我建议您在 Github 上提交一些内容。当我尝试这个问题时,一切都运行良好。你可以试试inspecting the dask graph manually。