【发布时间】:2022-10-15 01:48:23
【问题描述】:
我正在尝试增加一个计数器变量,该变量维护总操作计数,应该在并行进程之间共享以增加它,为此我得到的解决方案是 Ray 中的“Actor”,但它也不起作用。 a 的值没有增加,它只是增加 1 并保持不变。
似乎每个进程仍在创建自己的 Counter 对象副本。我怎样才能只用面向对象的方法做同样的事情?
当我使用 Python lib multiprocessing multiprocessing.Pool().apply_async(A.task1,callback=self.task2()) 时,同样的方法也有效。
我怎样才能在 Ray 中做同样的事情,或者在 Dask 中是否可行?
import ray, time
@ray.remote
class Counter:
def __init__(self):
self.a = 0
def inc_a(self):
self.a +=1
def get_a(self):
return self.a
class A:
def __init__(self) -> None:
self.b = 0
def dotask(self):
for _ in range(10):
# print(f"Before Counter(a: {ray.get(c_obj.a.remote())}, b: {self.b})")
c_obj = Counter.remote()
A.task1.remote(self, c_obj)
self.b += 1
# print(f"After Counter(a: {ray.get(c_obj.a.remote())}, b: {self.b})")
@staticmethod
@ray.remote
def task1(self, c_obj):
time.sleep(20)
self.task2(c_obj)
def task2(self, c_obj):
c_obj.inc_a.remote()
print(f"After Inc (a: {ray.get(c_obj.a.remote())}, b: {self.b})")
电流输出:
(1,0),(1,1),(1,2),(1,3),(1,4),(1,5),(1,6),(1,7),(1,8),(1,9)
预期输出:
(1,0),(2,1),(3,2),(4,3),(5,4),(6,5),(7,6),(8,7),(9,8),(10,9)
【问题讨论】:
标签: python python-multiprocessing dask ray