【问题标题】:Python Queue: how to delay / prolong / modify timeout for blocking get?Python 队列:如何延迟/延长/修改阻塞获取的超时?
【发布时间】:2020-03-17 04:55:50
【问题描述】:

我有一个由线程填充的queue.Queue。一个方法尝试从这个队列接收超时。现在假设,另一个线程可以重置我们队列等待的超时,如果队列没有及时送入,我们的接收函数应该继续,更新超时。 我可以按如下方式实现,但是我必须修改内置的queue.Queue 类,以便在等待期间可以修改get() 方法中的endtime 参数... 有没有更好的解决方案? (我不想使用 asyncio...)

from threading import Thread
from queue import Queue, Empty
import time

q = Queue()
TIMEOUT = 1
RESET_TIME = 0.5
PUT_TIME = 1.2
t0 = time.time()

def receive():
    try:
        _res = q.get(block=True, timeout=TIMEOUT)
        print(f'get @ {time.time()-t0}')
        return _res
    except Empty:
        print(f'to @ {time.time()-t0}')
        return None

def feed_queue():
    time.sleep(PUT_TIME)
    print(f'put @ {time.time()-t0}')
    q.put_nowait(42)

def reset_timeout():
    time.sleep(RESET_TIME)
    with q.mutex:
        q.endtime += TIMEOUT
    print(f'reset @ {time.time()-t0}')

if __name__ == '__main__':
    Thread(target=feed_queue).start()
    Thread(target=reset_timeout).start()
    res = receive()
    print('res:', res)

这会产生:

reset @ 0.5013222694396973
put @ 1.201164722442627
get @ 1.201164722442627
res: 42

在 queue.py 中进行了以下修改以使其正常工作:

Index: queue.py
===================================================================
--- queue.py    (revision 28725)
+++ queue.py    (working copy)
@@ -52,6 +52,7 @@
         # drops to zero; thread waiting to join() is notified to resume
         self.all_tasks_done = threading.Condition(self.mutex)
         self.unfinished_tasks = 0
+        self.endtime = 0

     def task_done(self):
         '''Indicate that a formerly enqueued task is complete.
@@ -171,9 +172,9 @@
             elif timeout < 0:
                 raise ValueError("'timeout' must be a non-negative number")
             else:
-                endtime = time() + timeout
+                self.endtime = time() + timeout
                 while not self._qsize():
-                    remaining = endtime - time()
+                    remaining = self.endtime - time()
                     if remaining <= 0.0:
                         raise Empty
                     self.not_empty.wait(remaining)

【问题讨论】:

    标签: python multithreading asynchronous queue timeout


    【解决方案1】:

    您可以创建自己的类,继承Queue 并添加全局变量endtime,如下所示:

    class Waszil(Queue):
        def __init__(self, maxsize=0):
            super().__init__(self)
            self.maxsize = maxsize
            self._init(maxsize)
            self.endtime = 0
    

    然后只需将 q = Queue() 更改为 q = Waszil() 就可以了。

    编辑: 如果你更喜欢在 Waszil 类中实现固有的线程安全,你可以像这样使用threading.Lock

    from threading import Lock
    
    class Waszil(Queue):
        def __init__(self, maxsize=0):
            super().__init__(self)
            self.threadLock = Lock()
            self.maxsize = maxsize
            self._init(maxsize)
            self.endtime = 0
    
        def increment_endtime(self):
            with self.threadLock:
                self.endtime += 1
    

    在这种情况下,而不是你的

    with q.mutex:
        q.endtime += TIMEOUT    
    

    您只需拨打q.increment_endtime()

    【讨论】:

    • 是的,我知道,但这仍然是相同的解决方案。我的问题是,如果有其他更惯用的方式来执行此操作,例如事件、条件等。我也不确定这种交互(更新此结束时间)是否是线程安全的,或者有任何流量。
    • with q.mutex: 应该确保线程安全。不过,我将在我的答案中添加一个EDIT 版本,它在Waszil 类中固有地处理线程安全。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-09-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多