如何使用Python Ray并行处理大量数据而不耗尽内存?

Joh*_*hsm 5 python-3.x ray

我正在考虑使用 Ray 来简单实现数据并行处理:

  • 有大量的数据项需要处理,这些数据项可以通过流/迭代器获得。每件物品的尺寸都很重要
  • 应该在每个项目上运行一个函数,并且会产生显着大小的结果
  • 处理后的数据应该在流中传递或存储在某种接收器中,该接收器只能在一段时间内接受一定量的数据

我想知道 Ray 是否可以做到这一点。

目前我有以下基于 python 多处理库的简单实现:

  • 一个进程读取流并将项目传递到队列,该队列将在 k 个项目后阻塞(以便队列所需的内存不会超过某个限制)
  • 有几个工作进程将从输入队列中读取并处理项目。处理后的项目被传递到结果队列,该队列的大小也是有限的
  • 另一个进程读取结果队列以传递项目

这样,一旦工作人员无法处理更多项目,队列就会阻塞,并且不会尝试将更多工作传递给工作人员。如果接收器进程无法存储更多项目,则结果队列将阻塞,这反过来又会阻塞工作线程,而工作线程又会阻塞输入队列,直到写入进程可以再次写入更多结果。

那么 Ray 有抽象来做这样的事情吗?我如何确保只有一定量的工作可以传递给工作人员,以及如何拥有类似单进程输出功能的东西,并确保工作人员不会用太多结果淹没该功能,以致内存/存储空间已耗尽?

小智 4

Ray 有一个实验性流 API,您可能会发现它很有用:https://github.com/ray-project/ray/tree/master/python/ray/experimental/streaming

它提供了流数据源、自定义运算符和接收器的基本构造。您还可以通过限制队列大小来设置应用程序的最大内存占用量。

您能否分享一些有关您的申请的其他信息?

我们正在谈论什么类型的数据?单个数据项有多大(以字节为单位)?