【问题标题】:How to use shared variable among parallel process and increment the counter如何在并行进程之间使用共享变量并增加计数器
【发布时间】: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


    【解决方案1】:

    您通过调用Counter.remote() 10 次来创建10 个不同的Counter 实例(参与者)。尝试将c_obj = Counter.remote() 移到for 循环之外,以便将相同的actor 传递给每个任务。

    如果您只想运行一批任务:

    @ray.remote
    def square(x):
        return x * x
    
    futures = [square.remote(i) for i in range(4)]
    
    print(ray.get(futures))
    # -> [0, 1, 4, 9]
    

    这是ray-core walkthrough 示例。


    另一个问题是在这里使用@staticmethod

    class A:
        @staticmethod
        @ray.remote
        def task1(self, c_obj):
    

    一个静止的方法没有得到一个特殊的self 成员,这意味着self 在这里只是一个普通的参数。因为这也是偏僻的可能在另一个节点上执行的方法,它的所有参数在调用它之前都被序列化,包括self

    因此,每次使用A.task1.remote(self, c_obj) 调用它时,实例self(A 类)的当前状态都会被序列化并传输到远程工作人员。这可能不是您想要的。

    【讨论】:

    • 感谢您提供的解决方案,它现在拥有相同的对象副本,但执行顺序不同..“a”的值没有像预期输出中显示的那样排序.. 是否可以让它像它工作的那样排序水池。应用异步。我可以在调用 A.task1.remote() 时设置 ray.get() 但这会减慢进程。在 apply_async 中它工作得很快。
    • 我不确定你现在的目标是什么。你真的想要一个共享的状态/计数器吗?还是只是为了让您可以使用它手动实现apply_async() 的替换?
    • 当前输出为 (4,0), (2,1) (7,2) (1,3) (5,4) (9,5) (6,6) (3,7) (0,8) ( 8,9) 和预期是 (1,0),(2,1),(3,2),(4,3),(5,4),(6,5),(7,6),( 8,7),(9,8),(10,9)
    • 您的预期输出将意味着任务没有并行运行。您正在从任务中调用 print 方法。计数器显示总操作计数(打印时)。
    猜你喜欢
    • 2021-03-27
    • 2012-04-09
    • 2013-12-28
    • 2012-06-17
    • 1970-01-01
    • 2016-12-14
    • 2012-03-24
    • 1970-01-01
    相关资源
    最近更新 更多