【问题标题】:Semaphores in dask.distributed?dask.distributed中的信号量?
【发布时间】:2018-07-17 22:53:20
【问题描述】:

我有一个包含 n 个工作人员的 dask 集群,并希望这些工作人员对数据库进行查询。但是数据库只能并行处理 m 个查询,其中 m

我已经看到分布式支持锁(http://distributed.readthedocs.io/en/latest/api.html#distributed.Lock)。但是这样一来,我只能并行执行一个查询,而不是 m。

我还看到我可以为每个工作人员定义资源 (https://distributed.readthedocs.io/en/latest/resources.html)。但这也不合适,因为数据库独立于工作人员。我要么必须为每个工作人员定义 1 个数据库资源(这会导致过多的并行查询)。或者我必须将 m 个数据库资源分配给 n 个工作人员,这在设置集群时很困难,并且在执行上也不是最优的。

是否可以在 dask 中定义信号量之类的东西来解决这个问题?

【问题讨论】:

    标签: dask dask-distributed


    【解决方案1】:

    您可以使用dask.distributed.Queue

    class DDSemaphore(object):
        """Dask Distributed Semaphore"""
    
        def __init__(self, value=1):
            self._q = dask.distributed.Queue()
            for _ in range(value):
                self._q.put(42)
    
        def acquire():
            self._q.get()
    
        def release():
            self._q.put(42)
    

    【讨论】:

      【解决方案2】:

      你可能会用锁和变量来破解一些东西。

      更简洁的解决方案是只实现信号量,就像实现锁的方式一样。根据您的经验,这可能并不难(锁实现是 150 行),并且会是一个受欢迎的拉取请求。

      https://github.com/dask/distributed/blob/master/distributed/lock.py

      【讨论】:

        猜你喜欢
        • 2018-05-27
        • 2017-08-24
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-03-25
        相关资源
        最近更新 更多