使用Queue()停止Python线程遇问题:错误时线程无法终止
问题分析与修复方案
核心问题点
- 读取线程未完成队列任务就退出:处理到"Read 7"时直接退出,但没对当前获取的队列项执行
q.task_done(),导致dlQueue的未完成任务计数未清零,主线程的dlQueue.join()会一直阻塞。 - 错误处理函数无效重置队列:
errorBreak里把全局队列重新赋值为空队列,但主线程和工作线程持有的还是原队列引用,这操作完全没用,还会导致状态混乱。 - 写入线程阻塞无法响应停止信号:写入线程先执行
q.get(),如果队列空会一直阻塞,根本没机会检查errorStop标志,自然无法退出。
修复后的代码
import pandas as pd import datetime import traceback from queue import Queue from threading import Thread import time dlQueue = Queue() writeQueue = Queue() dlQDone = False errorStop = False def log(text): text = datetime.datetime.now().strftime("%Y/%m/%d, %H:%M:%S ") + text print(text) def errorBreak(): global errorStop global dlQDone # 仅设置停止标志,无需重置队列 errorStop = True dlQDone = True def downloadTable(t, q): global dlQDone global errorStop while True: if errorStop: # 退出前尝试标记任务完成,避免队列阻塞 try: q.task_done() except ValueError: pass return try: # 带超时的队列获取,定期检查停止标志 nextQ = q.get(timeout=1) except: if errorStop: return continue log("READING: " + nextQ) writeQueue.put("Writing " + nextQ) log("DONE READING: " + nextQ) ####模拟错误场景### if nextQ == "Read 7": log("Breaking Read") errorBreak() q.task_done() # 必须标记当前任务完成 return ################### q.task_done() if q.qsize() == 0: log("Download QUEUE finished") dlQDone = True return def writeTable(t, q): global errorStop global dlQDone while True: if errorStop: log("Error Stop return") # 清空剩余队列并标记任务完成,确保join能结束 while not q.empty(): q.get() q.task_done() return try: # 带超时获取,避免永久阻塞 nextQ = q.get(timeout=1) except: if errorStop or (dlQDone and q.qsize() == 0): log("Writing QUEUE finished") return continue log("WRITING: " + nextQ) log("DONE WRITING: " + nextQ) q.task_done() if dlQDone and q.qsize() == 0: log("Writing QUEUE finished") return try: log("PROCESS STARTING!!") for i in range(10): dlQueue.put("Read " + str(i)) startTime = time.time() log("Starting threaded pull....") dlWorker = Thread(target=downloadTable, args=("DL", dlQueue,)) dlWorker.start() writeWorker = Thread(target=writeTable, args=("Write", writeQueue,)) writeWorker.start() dlQueue.join() writeQueue.join() log(f"Finished thread in {str(time.time() - startTime)} seconds") log(f"Threads: DL={dlWorker.is_alive()}, Write={writeWorker.is_alive()}") except Exception as error: log(error) log(traceback.format_exc())
关键修复说明
- 补全任务完成标记:读取线程错误退出前必须调用
q.task_done(),确保主线程join()能正常结束。 - 移除无效队列重置:
errorBreak仅设置停止标志,避免队列引用混乱。 - 超时获取队列项:读写线程都用带超时的
q.get(),避免永久阻塞,定期检查停止标志。 - 错误时清理队列:写入线程响应停止信号时,清空剩余队列并标记任务完成,保证
writeQueue.join()顺利结束。
内容的提问来源于stack exchange,提问作者Pearl
相关产品推荐
相关产品推荐

