Python多进程中能否同时使用Queue和Event?
Python多进程中能否同时使用Queue和Event?
当然可以同时用Queue和Event啦!你的问题其实不是这两个组件不能共存,而是worker的逻辑有个小漏洞——当你调用tasks.get()的时候,这个方法默认是阻塞式的:如果队列里没有任务,worker会一直卡在这行代码上,根本没机会去执行stop_flag.is_set()的检查,自然就接收不到停止信号了,主进程也就会一直卡在w.join()那里等着。
咱们来修改一下worker的逻辑,让它有机会定期检查停止信号,这里有个简单的解决办法:给tasks.get()加上超时时间,这样每隔一段时间,就算队列里没任务,worker也能从阻塞中醒过来,去检查stop_flag的状态。
修改后的代码示例:
import multiprocessing as mp from time import sleep from queue import Empty # 需要导入Empty异常 def worker(stop_flag: mp.Event, tasks: mp.Queue): while True: # 先检查停止信号,避免刚启动就阻塞在get上 if stop_flag.is_set(): print("Stop signal received") break try: # 带超时的get,超时后会抛出Empty异常 task = tasks.get(timeout=0.5) print(f"Processed task: {task}") # 如果是JoinableQueue的话,这里需要task_done() # tasks.task_done() except Empty: # 队列空的时候,继续循环检查停止信号 continue if __name__ == '__main__': tasks = mp.Queue() stop_flag = mp.Event() w = mp.Process(target=worker, args=(stop_flag, tasks)) w.start() sleep(1) items = 2 print("Starting work loop") for i in range(items): tasks.put(i) print(f"Put item #{i} of {items} in queue") sleep(2) stop_flag.set() print("Stop flag set") print("Stopping worker") w.join() print("Worker joined") w.close() print("Finished")
这样改了之后,worker就不会一直卡死在tasks.get()上了:当队列空的时候,每隔0.5秒它就会醒一次,去检查stop_flag有没有被设置,一旦检测到停止信号,就会退出循环。
另外还有个小细节:我把stop_flag.is_set()的检查放在了tasks.get()前面,这样如果主进程已经发送了停止信号,worker就不会再去尝试获取任务了,能更快响应停止指令。
如果你用的是JoinableQueue,记得在处理完任务后调用tasks.task_done(),不过核心逻辑还是一样的——要让worker有机会跳出阻塞,去检查停止信号。
备注:内容来源于stack exchange,提问作者Sergei Katkovsky
相关产品推荐
相关产品推荐

