如何在Python主线程中重新抛出无限运行IO线程的异常?
问题
我在平台上看到过许多类似问题,但现有答案都无法解决我的当前问题。
我开发了一个通过串行连接收发消息的程序,消息可随时接收,因此建立连接后会启动两个线程:一个从端口读取消息并加入队列,另一个处理收到的消息。这些线程需持续运行直至连接关闭。
默认情况下,Python处理线程异常的行为是仅将异常打印到终端并终止线程,这显然不符合我的需求,目前看来最佳解决方案是在主线程中重新抛出异常,以标识监听或处理线程因异常退出。
我该如何正确捕获程序中无限运行IO线程内的异常?
我尝试通过重写threading.excepthook自定义函数来重新抛出异常,但这并不奏效,因为该函数仍在抛出异常的线程内被调用,而非主线程。
以下是一段极简代码片段,大致展示当前代码逻辑:
import threading from time import sleep import random keep_thread_going = True def thread_func(delay): while keep_thread_going: sleep(delay) # 模拟随机抛出异常 if random.randint(0, 5) == 0: raise Exception("线程随机抛出了一个异常") def except_hook_test(args): print("自定义异常处理...") raise args.exc_value # 怎么才能把这个异常抛到主线程里? def main(): threading.excepthook = except_hook_test thread_obj = threading.Thread(target=thread_func, args=[1], daemon=True) thread_obj.start() # 模拟主线程做其他事情 sleep(5) # 终止线程 global keep_thread_going keep_thread_going = False thread_obj.join() print("完成") main()
解决方案
方法1:用共享异常队列
主线程维护一个线程安全的队列,子线程捕获异常后把它丢进队列,主线程定期检查队列,一旦有异常就直接抛出。
修改后的代码示例:
import threading from time import sleep import random from queue import Queue keep_thread_going = True exception_queue = Queue() def thread_func(delay): global keep_thread_going try: while keep_thread_going: sleep(delay) if random.randint(0, 5) == 0: raise Exception("线程随机抛出了一个异常") except Exception as e: exception_queue.put(e) # 异常发生后终止线程 keep_thread_going = False def main(): thread_obj = threading.Thread(target=thread_func, args=[1], daemon=True) thread_obj.start() # 主线程循环:做其他工作 + 检查异常队列 for _ in range(5): sleep(1) # 检查是否有异常 if not exception_queue.empty(): exc = exception_queue.get() raise exc # 如果没异常,正常终止线程 global keep_thread_going keep_thread_going = False thread_obj.join() print("完成") main()
方法2:自定义线程类存储异常状态
继承threading.Thread写个自定义线程类,重写run方法捕获异常并存在线程对象里,主线程随时检查这个对象,有异常就抛到主线程。
示例代码:
import threading from time import sleep import random keep_thread_going = True class ExceptionThread(threading.Thread): def __init__(self, target=None, args=(), kwargs=None): super().__init__(target=target, args=args, kwargs=kwargs) self.exc = None def run(self): try: super().run() except Exception as e: self.exc = e # 异常发生后终止线程 global keep_thread_going keep_thread_going = False def thread_func(delay): while keep_thread_going: sleep(delay) if random.randint(0, 5) == 0: raise Exception("线程随机抛出了一个异常") def main(): thread_obj = ExceptionThread(target=thread_func, args=[1], daemon=True) thread_obj.start() # 主线程做其他工作 for _ in range(5): sleep(1) # 检查线程是否有异常 if thread_obj.exc is not None: raise thread_obj.exc # 正常终止线程 global keep_thread_going keep_thread_going = False thread_obj.join() print("完成") main()
方法3:使用concurrent.futures.ThreadPoolExecutor
如果场景不复杂,直接用concurrent.futures.ThreadPoolExecutor,它的Future对象调用result()时会自动把子线程的异常抛到主线程里。
示例代码:
from concurrent.futures import ThreadPoolExecutor from time import sleep import random keep_thread_going = True def thread_func(delay): while keep_thread_going: sleep(delay) if random.randint(0, 5) == 0: raise Exception("线程随机抛出了一个异常") def main(): with ThreadPoolExecutor(max_workers=1) as executor: future = executor.submit(thread_func, 1) # 主线程做其他工作 for _ in range(5): sleep(1) # 检查future是否完成(异常会触发done) if future.done(): # result()会在主线程抛出子线程的异常 future.result() # 正常终止线程 global keep_thread_going keep_thread_going = False # 等待线程结束,若有异常也会抛出 future.result() print("完成") main()
内容的提问来源于stack exchange,提问作者hdconway
相关产品推荐
相关产品推荐

