You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何multiprocessing.pool.ThreadPool的terminate()会挂起?KeyboardInterrupt终止异常

解决KeyboardInterrupt终止异步多进程/线程任务时的挂起问题

我在尝试用KeyboardInterrupt终止异步多进程任务时,发现调用相关终止逻辑(比如设置stopEvent)后,程序有时会出现挂起现象。以下是我的相关代码:

from multiprocessing.pool import ThreadPool
import multiprocessing
import time
import queue
import inspect

def worker(index):
    print('{}: start'.format(index))
    for i in range(5):
        time.sleep(1)
    print('{}: stop'.format(index))
    return index, True

def wrapper(index, stopEvent, qResult):
    if stopEvent.is_set() is True:
        return index, False
    try:
        result = worker(index)
    except:
        print('*' * 50)
        return index, False
    else:
        if result[1] == True:
            qResult.put(result)
        return result

def watcher(qResult, stopEvent):
    cntQ = 0
    while True:
        try:
            result = qResult.get(timeout=10)
            qResult.task_done()
        except queue.Empty:
            if stopEvent.is_set() is True:
                break
        except KeyboardInterrupt:
            stopEvent.set()
        else:
            cntQ += 1
            print(result)
    qResult.join()
    qResult.close()
    print('qResult count:', cntQ)

def main():
    stopEvent = multiprocessing.Event()
    qResult = multiprocessing.JoinableQueue()
    qResult.cancel_join_thread()
    watch = multiprocessing.Process(target=watcher, args=(qResult, stopEvent))
    watch.start()
    pool = ThreadPool()
    lsRet = []
    for i in range(100000):
        try:
            ret = pool.apply_async(wrapper, args=(i, stopEvent, qResult))
            lsRet.append(ret)
        except KeyboardInterrupt:
            stopEvent.set()
            time.sleep(1)
            break
        if i+1 % 10 == 0:
            time.sleep(2)
    cntTotal = len(lsRet)
    cntRet = 0
    for ret in lsRet:
        if stopEvent.is_set():
            break
        try:
            ret.get()
        except KeyboardInterrupt:
            stopEvent.set()
            time.sleep(1)
        else:
            cntRet += 1
    # 省略后续清理代码

问题分析

你的代码里存在几个核心问题导致挂起:

  • 线程池清理不彻底:触发中断后仅设置stopEvent并跳出循环,但未主动终止线程池。工作线程可能仍在执行worker中的time.sleep或等待任务,无法正常退出。
  • 队列join()阻塞:watcher进程最后调用qResult.join(),若队列中还有未完成的任务(比如工作线程没来得及放入结果或调用task_done()),这个调用会一直阻塞,导致watcher进程无法退出。
  • 任务无法中途终止:wrapper仅在开头检查stopEvent,一旦进入worker函数,即使stopEvent被设置,任务也会继续执行到结束,拖慢终止流程。
  • 信号处理竞争:KeyboardInterrupt信号可能在多线程/进程间传递不及时,导致部分组件无法响应中断指令。

解决方案

针对这些问题,我们可以逐个修复,以下是优化后的完整代码:

from multiprocessing.pool import ThreadPool
import multiprocessing
import time
import queue
import signal

def worker(index, stopEvent):
    print('{}: start'.format(index))
    for i in range(5):
        # 循环中定期检查终止信号,支持中途退出
        if stopEvent.is_set():
            print('{}: interrupted'.format(index))
            return index, False
        time.sleep(1)
    print('{}: stop'.format(index))
    return index, True

def wrapper(index, stopEvent, qResult):
    if stopEvent.is_set():
        return index, False
    try:
        result = worker(index, stopEvent)
    except Exception as e:
        print(f'Worker {index} error: {str(e)}')
        return index, False
    else:
        if result[1]:
            try:
                qResult.put(result, timeout=1)
            except queue.Full:
                pass
        return result

def watcher(qResult, stopEvent):
    cntQ = 0
    # 让watcher进程忽略中断信号,由主进程统一处理
    signal.signal(signal.SIGINT, signal.SIG_IGN)
    while True:
        if stopEvent.is_set():
            # 清空队列,避免join阻塞
            while not qResult.empty():
                try:
                    qResult.get_nowait()
                    qResult.task_done()
                except queue.Empty:
                    break
            break
        try:
            result = qResult.get(timeout=1)
            qResult.task_done()
        except queue.Empty:
            continue
        else:
            cntQ += 1
            print(result)
    qResult.close()
    print('qResult count:', cntQ)

def main():
    stopEvent = multiprocessing.Event()
    qResult = multiprocessing.JoinableQueue()
    qResult.cancel_join_thread()
    watch = multiprocessing.Process(target=watcher, args=(qResult, stopEvent))
    watch.start()
    pool = ThreadPool()
    lsRet = []
    try:
        for i in range(100000):
            ret = pool.apply_async(wrapper, args=(i, stopEvent, qResult))
            lsRet.append(ret)
            if (i+1) % 10 == 0:
                time.sleep(2)
    except KeyboardInterrupt:
        print('Received KeyboardInterrupt, stopping all tasks...')
        stopEvent.set()
        # 立即终止线程池,停止所有未完成任务
        pool.terminate()
        pool.join()
    else:
        # 正常完成时关闭线程池
        pool.close()
        pool.join()
        stopEvent.set()
    
    # 等待watcher进程退出
    watch.join()
    print('All tasks stopped.')

if __name__ == '__main__':
    main()

关键修改说明

  1. worker支持中途终止:在循环中加入stopEvent检查,收到终止信号后立即退出任务,避免长时间阻塞。
  2. 统一中断处理:让watcher进程忽略SIGINT,由主进程统一处理中断,避免多进程信号竞争。
  3. 主动清理线程池:触发中断后调用pool.terminate()强制终止所有工作线程,再用pool.join()等待清理完成。
  4. 队列安全退出:watcher退出前清空队列,避免join()因未完成任务阻塞。
  5. 优化异常信息:捕获异常时输出具体错误,方便排查问题。

这样修改后,程序在收到KeyboardInterrupt时能快速、可靠地终止所有任务,避免挂起现象。

内容的提问来源于stack exchange,提问作者snoman

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 03:51:13