【问题标题】:How to put an item back to a queue.Queue如何将项目放回队列。队列
【发布时间】:2020-10-16 15:39:16
【问题描述】:

如何将项目返回到 queue.Queue?如果任务失败,这在线程或多处理中很有用,因此任务不会丢失。

docs for queue.Queue.get() 说函数可以“从队列中删除并返回一个项目”,但我相信这里使用的“返回”一词是指函数将项目返回给调用线程,而不是放置它回到项目队列。下面的示例代码证明了这一点,它只是在主线程的第二个 queue.Queue.get() 调用上无限阻塞,而不是在线程中的 print() 调用上进行。

import time
import threading
import queue


def threaded_func():
    thread_task = myqueue.get()
    print('thread_task: ' + thread_task)

myqueue = queue.Queue()
myqueue.put('some kind of task')
main_task = myqueue.get()
print('main_task: ' + main_task)

t = threading.Thread(target=threaded_func)
t.daemon = True
t.start()

time.sleep(5)
myqueue.get()   # This blocks indefinitely

我必须相信有一种简单的方法可以将任务放回去,那是什么?调用task_done(),然后调用put() 并在两个操作中将其放回队列中的任务不是原子的,因此可能会导致丢失项目。

一种可能但笨拙的解决方案是尝试再次执行任务,但是您必须添加一些额外的行来处理这种复杂性,我什至不确定所有失败的任务是否一定会以这种方式恢复。

【问题讨论】:

  • 您在解释问题,但并没有真正解释总体目标。那么你到底想在这里完成什么?
  • 这是为了不丢失失败的工作项,以便以后能够处理它。我的具体用途是调用 REST API URL,当工作项 URL 不固有的一些问题可能导致它失败时。

标签: python multithreading queue


【解决方案1】:

并非所有失败的任务都可以恢复。你不应该重试它们,除非有某种理由认为它们会在以后通过。例如,如果您的工作项是一个 URL 并且连接失败计数,您可以实现某种最大重试次数。

你最大的问题是你还没有实现一个可行的工人模型。您需要 2 个队列才能与工作人员进行双向对话。一个用于发布工作项目,一个用于接收状态。一旦你有了它,接收者总是可以决定将该消息塞回工作队列中。这是一个懒惰的工人的例子,它只是通过它告诉的内容。

import threading
import queue

def worker(in_q, out_q):
    while True:
        try:
            task, data = in_q.get()
            print('worker', task, data)
            if task == "done":
                return
            elif task == "pass this":
                out_q.put(("pass", data))
            else:
                out_q.put(("fail", data))
        except Exception as e:
            print('worker exception', e)
            out_q.put("exception", data)

in_que = queue.Queue()
out_que = queue.Queue()

work_thread = threading.Thread(target=worker, args=(in_que, out_que))
work_thread.start()

# lets make every other task a fail
in_que.put(('pass this', 0))
in_que.put(('fail this', 1))
in_que.put(('pass this', 2))
in_que.put(('fail this', 3))
in_que.put(('pass this', 4))
in_que.put(('fail this', 5))

pending_tasks = 6

while pending_tasks:
    status, data = out_que.get()
    if status == "pass":
        pending_tasks -= 1
    else:
        # make failing tast pass
        in_que.put(('pass this', data))

in_que.put(("done", None))
work_thread.join()
print('done')

【讨论】:

  • 我的演示代码只是为了说明无法将项目放回队列,它不是我的实际工作代码。
  • 哎呀,我不知道键盘输入会创建帖子而不是换行符。我还想说,我可以看到您的代码正在使用第二个队列,并对第一个队列中的项目进行计数,如果它们失败则将它们重新排队。我认为这个额外的队列和处理逻辑是一种不必要的复杂性,当人们可以将一个项目返回主队列时,这就是我最初的问题所在。那么您是说不能将项目返回到队列中?
  • 嗯,当然。我所做的只是将消息放回主队列。另一个队列用于处理线程的返回。你的代码不可能工作......进入主线程应该做什么,究竟是什么?所以我写了代码。 “不能工作”到“工作”也许是最小的复杂性。
【解决方案2】:

您可以使用PriorityQueue,其中条目按优先顺序返回。

from tornado.queues import PriorityQueue

q = PriorityQueue()
q.put((2, 'item 1'))
q.put((2, 'item 2'))

q.get()  # item 1
q.put((1, 'item 1'))  # put item 1 back on the queue

q.get()  # item 1
q.get()  # item 2

请注意,get() 的返回也是一个优先级编号和项目的元组。

【讨论】:

    猜你喜欢
    • 2021-12-23
    • 2012-03-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-07-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多