【问题标题】:gevent - Pass a value to many greenlets in parallelgevent - 将一个值并行传递给许多 greenlets
【发布时间】:2014-10-02 14:34:09
【问题描述】:

我正在尝试实现一个简单的 gevent 设置。有一个发送者应该并行向几个服务员发送一个值。 Event 类最接近于解决这个问题,如下所示。

每三秒,setter 会创建一个事件,该事件会解除对所有服务员的阻塞。事件随即被清除,因此服务员再次阻塞,直到下一次。

import gevent
from gevent.event import Event

evt = Event()

def setter():
    '''After 3 seconds, wake all threads waiting on the value of evt'''
    while True:
        gevent.sleep(3)
        evt.set()
        evt.clear()

def waiter(arg):
    while True:
        evt.wait()
        print("waiter {}".format(arg))

def main():
    gevent.joinall([
        gevent.spawn(setter),
        gevent.spawn(waiter,1),
        gevent.spawn(waiter,2),
        gevent.spawn(waiter,3),
    ])

if __name__ == '__main__': main()

现在,我只需要执行此操作,还需要向服务员传递一个值。显而易见的选择是使用AsyncResult。但是,无法clear AsyncResult 对象,因此等待者最终陷入无限循环。

你有什么想法如何实现这个吗?

【问题讨论】:

  • 为什么不使用共享变量来存储您想要发送的任何内容,并将引用传递给设置器和服务员。 setter 可以在调用set之前将其更改为任何内容?

标签: python concurrency gevent


【解决方案1】:

为什么不直接使用可变对象来存储您想要发送的任何内容,并将引用传递给 setter 和 waiters。 setter 可以在调用 set 之前将其更改为任何内容?见下文:

import gevent
from gevent.event import Event

evt = Event()

def setter(arg):
    '''After 3 seconds, wake all threads waiting on the value of arg['it']'''
    while True:
        gevent.sleep(3)
        arg['it'] += 1
        evt.set()
        evt.clear()

def waiter(num, arg):
    while True:
        evt.wait()
        print("waiter {} {}".format(num, arg['it']))

def main():
    THING = {'it': 1}
    gevent.joinall([
        gevent.spawn(setter, THING),
        gevent.spawn(waiter, 1, THING),
        gevent.spawn(waiter, 2, THING),
        gevent.spawn(waiter, 3, THING),
    ])

if __name__ == '__main__': main()

输出:

waiter 2 2
waiter 1 2
waiter 3 2
waiter 2 3
waiter 1 3
waiter 3 3
waiter 2 4
waiter 1 4
waiter 3 4

【讨论】:

    【解决方案2】:

    我认为最好的办法是为此使用队列。我创建了一个BroadcastQueue 类,它可以更轻松地管理向许多消费者发送一个值。生产者调用BroadcastQueue.broadcast(),它将向所有注册的消费者发送一个值。消费者通过调用BroadcastQueue.register 进行注册,这将返回一个唯一的gevent.queue.Queue() 对象。然后,消费者使用该对象来get 来自生产者的消息。

    import gevent
    from gevent.queue import Queue
    
    
    class BroadcastQueue(object):
        def __init__(self):
            self._queues = []
    
        def register(self):
            q = Queue()
            self._queues.append(q)
            return q
    
        def broadcast(self, val):
            for q in self._queues:
                q.put(val)
    
    
    def setter(bqueue):
        '''After 3 seconds, wake all threads waiting on the value of evt'''
        while True:
            gevent.sleep(3)
            bqueue.broadcast("hi")
    
    def waiter(arg, bqueue):
        queue = bqueue.register()
        while True:
            val = queue.get()
            print("waiter {} {}".format(arg, val))
    
    def main():
        bqueue = BroadcastQueue()
        gevent.joinall([
            gevent.spawn(setter, bqueue),
            gevent.spawn(waiter, 1, bqueue),
            gevent.spawn(waiter, 2, bqueue),
            gevent.spawn(waiter, 3, bqueue),
        ])
    
    if __name__ == '__main__':
        main()
    

    输出:

    waiter 1 hi
    waiter 2 hi
    waiter 3 hi
    waiter 1 hi
    waiter 2 hi
    waiter 3 hi
    waiter 1 hi
    waiter 2 hi
    waiter 3 hi
    waiter 1 hi
    waiter 2 hi
    waiter 3 hi
    waiter 1 hi
    waiter 2 hi
    waiter 3 hi
    

    【讨论】:

    • 谢谢!这非常有用。我觉得这应该默认包含在gevent 中。他们在编写文档方面已经很糟糕了,提供更全面的工具包不会有什么坏处..
    猜你喜欢
    • 2014-02-01
    • 2021-10-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-07-30
    相关资源
    最近更新 更多