焦小花同学

Ray,分布式Python让您的应用轻松扩展!

点击

上方文字

关注我们

今天我们聊聊一个超神奇的库——Ray。Ray 是一个分布式计算框架,它让你可以轻松地在多核、多机器上并行处理任务,甚至不用关心底层的复杂性。是的,你没听错,Ray 可以帮你把代码无缝扩展到集群上,而你几乎不需要改变现有的代码。适合那些想要让应用跑得更快、处理更多数据的小伙伴。

1. 

什么是 Ray?

Ray 是一个开源的 Python 库,专门为分布式计算设计。它允许你写出并行和分布式应用,而不用操心许多底层的分布式系统细节。Ray 的核心功能包括任务并行、Actor 模型、集群管理等。简单来说,它让你可以在多核 CPU 上并行化任务,甚至可以跨多台机器分布式执行任务。

温馨提示:

Ray 是开源的,但它非常强大且灵活。如果你只是玩玩,那用 Ray 的本地模式就够了,但如果你要处理大规模的数据,Ray 也可以扩展到多个节点。

2. 

Ray 的基础用法

我们先看看 Ray 是怎么工作的。Ray 的核心概念包括任务 (Task) 和 Actor 模型。任务是并行执行的函数,Actor 则可以理解为有状态的并行执行单元,可以长时间运行。

Task:并行处理的基础

Ray 让你可以把函数变成并行任务。好啦,直接上代码吧!

1import ray

2import time

4# 初始化 Ray

5ray.init()

7@ray.remote

8def my_task(x):

9 time.sleep(1)

10 return x * 2

12# 调用任务

13result = my_task.remote(5)

15# 获取结果

16print(ray.get(result)) # 输出 10

这段代码里,我们用 @ray.remote 把函数 my_task 变成了一个可以并行执行的任务。我们用 .remote() 方法调用它,而不是直接调用函数。这个 .remote() 会立即返回一个未来对象,Ray 会在后台并行执行任务。而 ray.get(result) 是用来获取任务的结果的。

温馨提示:

很多人会忽略 .remote() 这个方法,直接调用函数,这样的话就不会触发 Ray 的并行机制哦。

3. 

Actor:有状态的并行

如果你的函数需要保存一些状态,比如一个计数器,或者更复杂的东西,Ray 提供了 Actor 模型。这些 Actor 是有状态的,并且可以长时间运行。让我们看看例子:

1@ray.remote

2class Counter:

3 def __init__(self):

4 self.count = 0

6 def increment(self):

7 self.count += 1

8 return self.count

10# 创建一个 Actor

11counter = Counter.remote()

13# 调用 Actor 方法

14print(ray.get(counter.increment.remote())) # 输出 1

15print(ray.get(counter.increment.remote())) # 输出 2

Counter 类通过 @ray.remote 修饰变成了一个 Actor。我们创建了一个远程的 counter 实例,并通过 .remote() 方法调用它的 increment 方法。这就是 Actor 的强大之处:它可以保持状态并在多个任务调用中持续运行。

温馨提示:

Actor 的方法调用也是并行的,但它们会按顺序执行,确保状态的一致性。

4. 

实战:分布式数据处理

了解了 Ray 的基础概念后,我们来看看它在实际场景中的应用吧。比如说,我们有一大堆数据需要处理,单机处理太慢了,而且数据量大到内存都不够用了。这时候,Ray 就能帮我们把任务分发到多个机器上。

我们来模拟一个简单的场景:处理一大堆数字,把每个数都平方,然后求和。

1import ray

2import numpy as np

4# 初始化 Ray

5ray.init()

7@ray.remote

8def process_chunk(chunk):

9 return np.sum(chunk ** 2)

11# 分割数据

12data = np.random.rand(10000000)

13chunks = np.array_split(data, 10)

15# 并行处理每个数据块

16results = [process_chunk.remote(chunk) for chunk in chunks]

18# 收集结果

19total_sum = sum(ray.get(results))

21print(f“Total sum: {total_sum}”)

这个例子里,我们把数据分成了 10 份,然后用 Ray 并行处理每一块数据,最后把结果汇总。想象一下,如果你有 10 台机器,这个任务会变得更快。

温馨提示:

ray.get() 会阻塞,直到所有的任务完成。如果你只想处理一部分结果,可以使用 ray.wait()。

5. 

Ray 的扩展性

Ray 的一个大招就是,它不仅支持本地的多核并行计算,还支持集群模式。也就是说,你可以轻松地把任务分发到整个集群上,而不需要修改代码。只要在不同的机器上安装好 Ray,启动集群,就可以直接用。

1ray start --head --port=6379

在集群的其他节点上:

1ray start --address='IP地址:6379'

这样,Ray 就会把你的任务分发到不同的机器上,做到真正的分布式计算。

温馨提示:

如果你要在生产环境中使用 Ray 的集群模式,最好对网络和安全进行充分的配置。

6. 

结尾

Ray 的核心思想就是让你不用操心底层的分布式复杂性,轻松地并行、扩展你的应用。它适合从本地的小项目到大规模的集群计算。Ray 的 API 也非常简洁,几乎没有学习门槛,你可以很快上手。

Ray 不仅仅是一个并行计算的工具,它还有很多其他强大的功能,比如分布式训练(用于机器学习)、分布式 Reinforcement Learning 等等,值得深入挖掘。

别忘了点

Image

分享、

Image

收藏、

Image

在看、

Image

点赞

哦!