登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  python教程

Python multiprocessing共享状态与进程安全队列的选择

来源:17golang原创

时间:2026-09-24 17:37:57 316浏览 收藏

Python 的 multiprocessing 里,Queue 和 Manager 都能让进程“互相看到东西”,但它们解决的不是同一个问题:Queue 传递消息或任务,接收方拿到的是序列化后重建的对象;Manager 维护一份由服务器进程托管的共享对象,其他进程通过代理访问。生产任务分发优先选 Queue,只有确实需要多个进程共同读写一份小型状态时才考虑 Manager。

要点速览
  • 任务、结果、事件通知适合用 Queue,数据边界清晰,生产者和消费者解耦。
  • 计数器、少量配置、共享字典适合用 Manager,但每次代理读写都有跨进程通信成本。
  • 大对象、热点写入和高频共享状态不要靠 Manager 硬撑,应改成消息聚合、批量提交或专门的共享内存方案。

先区分消息传递和共享状态

我在拆分批处理脚本时,最容易犯的错误是把“结果要汇总”理解成“所有进程都要共享一个列表”。如果 worker 只是把完成项交回主进程,Queue 已经足够;如果多个进程必须读取并更新同一个小字典,Manager 才有明确价值。

这条边界也解释了为什么普通的全局变量不行:进程有各自的地址空间,子进程修改自己的副本,不会自动改到父进程。Queue 的代价是对象会经过 pickle,拿到的是副本;Manager 的代价是代理方法要访问管理器进程,不能当成本地字典无限次调用。

Queue适合把任务交给进程,而不是共享可变对象

下面的结构把任务输入和结果输出分成两条消息通道。它适合一个主进程派发多个独立任务,worker 只处理消息并返回结果。示例刻意用小字典传递业务数据,不依赖共享列表。

from multiprocessing import Process, Queue

def worker(tasks, results):
    # worker 只消费消息,不直接修改主进程里的容器
    while True:
        job = tasks.get()
        if job is None:
            break  # None 是本例约定的停止标记
        try:
            value = job["value"] * 2
            results.put({"id": job["id"], "value": value})
        except (KeyError, TypeError) as exc:
            # 把可预期的输入错误作为结果返回,避免静默丢任务
            results.put({"id": job.get("id"), "error": str(exc)})

if __name__ == "__main__":
    tasks, results = Queue(), Queue()
    workers = [Process(target=worker, args=(tasks, results)) for _ in range(2)]
    for process in workers:
        process.start()
    for job_id, value in enumerate([3, 5, 8], 1):
        tasks.put({"id": job_id, "value": value})
    for _ in workers:
        tasks.put(None)  # 每个 worker 都需要一个停止标记
    for _ in range(3):
        print(results.get(timeout=5))  # 读取结果时设置超时,避免永久等待
    for process in workers:
        process.join()  # 主进程明确回收子进程
Python multiprocessing Queue 将任务字典序列化后交给 worker 并由结果队列返回结果的结构说明图
图1:Queue 消息边界说明图,展示任务队列、worker 和结果队列之间的静态关系,不是运行截图。

这里不要用 empty() 判断“已经没有结果”,因为生产者和队列底层 feeder thread 之间可能存在短暂间隔。更稳妥的做法是根据任务数量收集结果,或用明确的结束消息和超时策略。

Manager适合共享小型状态,但代理调用要收敛

Manager() 会启动一个管理器进程,dict()list() 等返回的是代理对象。它适合共享少量状态,例如每个 worker 的完成计数或一份小配置。不要在循环里对代理对象做数十万次细粒度更新;可以让 worker 本地聚合,最后一次性提交。

from multiprocessing import Manager, Process

def worker(worker_id, shared):
    # 先在本地累加,减少对 Manager 代理的往返调用
    local_count = 0
    for _ in range(100):
        local_count += 1
    shared[worker_id] = local_count  # 最后只写入一次共享状态

if __name__ == "__main__":
    with Manager() as manager:
        progress = manager.dict()
        workers = [Process(target=worker, args=(i, progress)) for i in range(2)]
        for process in workers:
            process.start()
        for process in workers:
            process.join()  # 先确认 worker 结束,再读取最终快照
        snapshot = dict(progress)  # 尽早转成本地副本,后续计算不再走代理
        print(snapshot)
Python multiprocessing Manager 服务器进程通过 dict 代理为多个 worker 提供共享小型状态的结构说明图
图2:Manager 代理边界说明图,展示管理器服务器、代理容器和进程访问关系,不是运行截图。

还有一个常见坑:把普通嵌套字典放进 Manager 列表后,直接修改嵌套字典通常不会自动触发代理同步。需要把嵌套容器也做成 Manager 代理,或者取出对象、修改后重新赋回。共享状态越复杂,越应该回到 Queue 的事件/快照模型。

按数据量和一致性做选择

场景优先方案原因与边界
任务分发、结果回收QueueFIFO、多生产者多消费者,消息所有权清楚
少量计数、状态快照Manager.dict/list写法直观,但代理访问有通信成本
大块数组或高频数值读写shared_memory 等专用方案避免把大量对象反复 pickle 或代理调用
跨机器共享外部队列或存储服务Manager 可以远程访问,但认证、可用性和运维边界更重

最后检查三件事:是否给每个 worker 发送了停止信号;是否对 get()put() 设置了合适的超时;是否对启动的进程逐个 join()。如果答案是否定的,先补齐生命周期,再讨论 Queue 还是 Manager。

常见问题

Queue 里的对象修改后会同步回原对象吗?

不会。对象会被序列化,接收端拿到的是重建后的副本;要传回修改结果,需要再次放入结果队列。

Manager 一定比 Queue 慢吗?

不能简单下结论,但代理访问要经过管理器进程,细粒度高频读写通常更贵。先本地聚合,再批量提交,通常更容易控制成本。

什么时候不该使用 Manager?

共享对象很大、更新频繁、需要高吞吐,或状态可以通过事件重放得到时,不要把 Manager 当数据库使用,优先采用 Queue、共享内存或专门存储。

声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>