首页
/ FaceSwap 队列管理器(lib/queue_manager):多进程流水线中的 EventQueue 与统一关机机制

FaceSwap 队列管理器(lib/queue_manager):多进程流水线中的 EventQueue 与统一关机机制

2026-09-05 13:29:32作者:侯霆垣

FaceSwap 的 convert(换脸转换)流程是一条由「读盘 → 推理检测 → 多线程 patch → 写盘」组成的流水线,各环节之间通过队列解耦。本文围绕 docs/full/lib/queue_manager.rst 对应的核心模块 lib/queue_manager.py 展开,讲清楚 EventQueue_QueueManager 的设计意图、各方法参数语义,以及它们在 convert、preview 等真实流水线中的调用方式和基于 EOF 的优雅停机协议。读完后,你将能够理解 FaceSwap 如何用「全局 shutdown 事件 + 哨兵消息」让所有工作线程/进程在出错时可靠退出,并知道如何用自己的方式接入这套队列体系。

为什么队列管理要单独放在一个模块里

打开 lib/queue_manager.py,文件头部就有一条醒目的注释(第 2–5 行):

""" Queue Manager for faceswap

    NB: Keep this in it's own module! If it gets loaded from
    a multiprocess on a Windows System it will break Faceswap"""

从源码结构看,这条注释揭示了一个 Windows 平台上的多进程陷阱:multiprocessing 模块在 spawn 子进程时会重新导入若干模块,而队列对象内部依赖的 threading.Event/锁等对象在子进程中若被重新初始化会导致行为异常。因此 FaceSwap 把队列管理器隔离在一个独立模块中,确保所有进程共享同一个模块级单例 queue_manager(定义于 lib/queue_manager.pyqueue_manager = _QueueManager())。这也是理解后文所有调用方式的前提:不要直接实例化 _QueueManager,而要通过模块导入拿到全局单例

模块顶部的导入也交代了它的技术底座:

  • from queue import Queue, Empty as QueueEmpty:基于标准库 queue.Queue 构建跨线程安全队列;
  • from lib.utils import get_module_objects:模块末尾通过 __all__ = get_module_objects(__name__) 动态生成导出列表(lib/queue_manager.py),这是 FaceSwap 全项目统一的 __all__ 生成方式。

EventQueue:带全局关机信号的队列

EventQueue 是对标准库 Queue 的薄封装(lib/queue_manager.py):

class EventQueue(Queue):
    def __init__(self, shutdown_event: threading.Event, maxsize: int = 0) -> None:
        super().__init__(maxsize=maxsize)
        self._shutdown = shutdown_event

    @property
    def shutdown_event(self) -> threading.Event:
        return self._shutdown

参数说明(继承自源码 docstring):

参数 类型 说明
shutdown_event threading.Event 所有受管队列共享的全局关机事件。一旦置位,表示主进程(以及这些队列)应当关闭
maxsize int,可选 队列可容纳条目的上限,0 表示不限长。默认 0

关键点在于:每个队列自身只负责传递数据,"关机"这件事被上提为队列之外的一个共享 threading.Event。工作线程在阻塞 get() 拿不到数据时,可以先轮询 queue.shutdown_event.is_set(),从而不必依赖"数据流耗尽"来判断是否该退出——这在出错、用户中断等场景下尤为重要。

_QueueManager:命名队列的注册表

_QueueManager 维护一个 {队列名: EventQueue} 字典和一个全局 threading.Eventlib/queue_manager.py)。其 docstring 明确提示:

Don't import this class directly, instead import via :func:queue_manager

它提供的公开方法及其参数语义如下。

add_queue(name, maxsize=0, create_new=False)

向管理器注册一个 EventQueue,返回最终生成的队列名(str):

  • name(str):队列名称;
  • maxsize(int,可选):队列上限,0 为不限长,默认 0
  • create_new(bool,可选):重名策略开关。
    • False(默认):若同名队列已存在,抛出 ValueError: Queue '<name>' already exists.,防止误创建重复队列;
    • True:若同名队列已存在,则在原名后追加递增整数(name0name1……)直到名称唯一,并返回改写后的名称(lib/queue_manager.py)。

注意去重逻辑的写法是 while name in self.queues: name = f"{name}{i}"——从源码结构看,i 始终为 0,即追加的只是固定的 "0" 后缀;只要新名字不冲突即可退出循环。

get_queue(name, maxsize=0)

"获取或创建"语义(lib/queue_manager.py):若队列存在则直接返回;不存在则以 maxsize 创建后返回。maxsize 仅在队列尚未存在时生效。这是实际代码中最常用的入口,因为它天然幂等,多线程同时启动时各线程调用也不会互相冲突。

del_queue(name)

按名删除队列,队列必须已存在于管理器中,否则直接触发 KeyErrorlib/queue_manager.py)。

terminate_queues()

全量终止(lib/queue_manager.py),执行顺序为:

  1. self.shutdown.set() —— 置位全局关机事件,所有持有 EventQueue.shutdown_event 的监听方都会看到关机信号;
  2. _flush_queues() —— 清空所有队列中堆积的数据;
  3. 对每个队列 queue.put("EOF") —— 向每个队列写入字符串哨兵 "EOF"

docstring 说明其用途是"存在错误时调用"。可以看到这里把 EOF 字符串正式确立为 FaceSwap 流水线内部约定的停机哨兵值

flush_queue(name) / _flush_queues()

flush_queue 循环 queue.get(True, 1) 直到队列为空(lib/queue_manager.py);私有方法 _flush_queues 遍历所有队列逐个冲刷,供 terminate_queues 内部使用。

debug_monitor(update_interval=2)

调试工具(lib/queue_manager.py):启动一个守护线程,每 update_interval 秒(默认 2 秒)按名称排序打印所有队列的 qsize()。在 scripts/train.py 中可以看到它被注释掉的调用示例:# from lib.queue_manager import queue_manager; queue_manager.debug_monitor(1)——开发时取消注释即可观测流水线各段积压情况,排查"卡在哪一级队列"非常直观。

真实调用链:convert 流程中的三个队列

下面以 FaceSwap 最核心的换脸脚本 convert 为例,说明这套机制如何被使用(以下均来自 scripts/convert.py)。

队列注册与容量策略

Convert.process 初始化阶段调用 _add_queues()scripts/convert.py):

def _add_queues(self) -> None:
    """Add the queues for in, patch and out."""
    for qname in ("convert_in", "convert_out", "patch"):
        queue_manager.add_queue(qname, self._queue_size)

三条队列分别对应流水线三级:convert_in(送入转换)、convert_out(patch 结果输出给写盘线程)、patch(待 patch 任务进入多线程池)。队列容量由 _queue_size 属性决定(scripts/convert.py):单进程模式(-sp)或 jobs == 1 时为 2,否则为 4——这是典型的有界队列背压控制,避免推理端跑太快把写盘端内存打爆。

工作线程与队列的接线

_get_threads()scripts/convert.py)把队列接到 MultiThread 上:

save_queue = queue_manager.get_queue("convert_out")
patch_queue = queue_manager.get_queue("patch")
return MultiThread(self._converter.process, patch_queue, save_queue,
                   thread_count=self._pool_processes, name="patch")

这里用的是 get_queue 而非 add_queue——注册已由 _add_queues 完成,这里只取引用。self._pool_processes 依据 -sp / -j 参数与 CPU 总数、图片数量取最小值(scripts/convert.py),与有界队列配合形成"进程数 × 队列深度"的内存预算。

EOF 哨兵:优雅停机协议

正常结束与异常结束共用同一套哨兵机制:

  • 正常收尾_convert_images() 在 patch 线程池和写盘线程都完成后,显式向输出队列塞入哨兵(scripts/convert.py):queue_manager.get_queue("convert_out").put("EOF")
  • 异常收尾process() 的 try 块在成功后调用 queue_manager.terminate_queues()scripts/convert.py);而程序级统一兜底在 lib/utils.pysafe_shutdown() 中——它在导入单例后调用 queue_manager.terminate_queues()sys.exitlib/utils.py),无论错误从何处抛出,消费端都能收到 EOF

消费端如何识别哨兵?看 lib/convert.py 中 patch 工作循环(lib/convert.py):

while True:
    inbound: T.Literal["EOF"] | ConvertItem | list[ConvertItem] = in_queue.get()
    if inbound == "EOF":
        logger.debug("EOF Received")
        # Signal EOF to other processes in pool
        in_queue.put(inbound)   # 把 EOF 放回队列,通知池内其他工作进程
        break

注意类型标注 T.Literal["EOF"] | ConvertItem | list[ConvertItem] 直接声明了队列负载的三种合法取值,以及一个细节:收到 EOF 后会先把它放回队列再 break,这样同一队列的其余并行 worker 也能取到哨兵并相继退出——这是一个"广播式"停机技巧。此外 patch 循环内部对单张图片的异常做了降级处理:转换失败时输出原图并继续(lib/convert.py),保证个别坏帧不会中断整条流水线。

preview 工具的同构用法

tools/preview/preview.py 复用了完全相同的模式:预测端从 queue_manager.get_queue("preview_predict_in") 取帧(tools/preview/preview.py),patch 阶段使用 preview_patch_in / preview_patch_out 两条队列。这说明命名约定(<模块>_in / <模块>_out)+ 全局单例是 FaceSwap 各工具间可复用的流水线基础设施;而 lib/convert.py 顶部仅在 TYPE_CHECKING 分支导入 from lib.queue_manager import EventQueue 用于类型标注(lib/convert.py),运行期通过队列对象传递,不直接依赖管理器类本身。

使用小结

场景 推荐做法
首次注册队列 queue_manager.add_queue(name, maxsize),重名直接报错更安全
并行代码取队列引用 queue_manager.get_queue(name),幂等无冲突
判断是否应退出 阻塞取数与 queue.shutdown_event.is_set() 轮询结合,或直接约定哨兵值
出错/结束 queue_manager.terminate_queues()(置位事件 + 冲刷 + 广播 EOF)
调试验积压 queue_manager.debug_monitor(interval) 守护线程周期打印 qsize

从源码结构看,lib/queue_manager.py 全文不足 200 行,职责单一:它不做生产消费,只负责"命名注册 + 共享关机事件 + EOF 广播"三件事。真正复杂的编排逻辑(有界容量选择、多进程池、线程错误监控)分散在 scripts/convert.pytools/preview/preview.py 等使用方中。阅读或扩展 FaceSwap 的多进程功能时,建议先确认目标模块是否可被子进程安全导入——这正是该模块头部那条 Windows 注释要提醒开发者的设计约束。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388