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

基于multiprocessing的多进程通信及worker调度实现咨询

多进程场景下跨进程通信及Worker调度实现方案

核心方案思路

使用Python标准库multiprocessing自带的进程安全队列作为通信载体,Handler作为全局调度中心持有公共队列,所有子进程启动时传入队列实例,即可实现子进程向父进程上报、父进程向子进程下发指令、子进程间通信的需求。

具体实现步骤

1. 通信载体定义

在Handler初始化阶段创建两类核心队列:

  • 全局上报队列:所有子进程(IMAP监控进程、账号创建Worker进程)均可向该队列写入消息,Handler进程循环消费队列中的消息,执行对应的调度逻辑
  • 任务下发队列:Handler将需要执行的任务写入该队列,所有Worker进程共同消费该队列抢单执行
    如果需要实现指定子进程的点对点通信,可以为每个子进程单独创建专属队列,父进程持有所有子进程的专属队列写入端,子进程持有自己专属队列的读取端即可。

2. 子进程改造

不管是IMAP类还是Worker类,都改造为继承multiprocessing.Process,初始化时传入需要的队列实例:

IMAP进程改造示例

import multiprocessing
import configparser
import imaplib

class IMAPProcess(multiprocessing.Process):
    # 初始化时传入全局上报队列
    def __init__(self, report_queue):
        super().__init__()
        self.report_queue = report_queue
        conf = configparser.ConfigParser()
        conf.read('config.ini')
        self._config = conf['imap']
        self._mail = imaplib.IMAP4_SSL(self._config['host'])
        self._mail.login(self._config['username'], self._config['password'])
        self._mail.list()

    def run(self):
        while True:
            # 原有监控邮箱逻辑
            new_mail = self._check_new_mail() # 替换为你自己的邮件检测逻辑
            if new_mail:
                # 解析邮件得到要执行的操作,构造消息写入上报队列
                msg = {
                    "type": "new_mail",
                    "action": "xxx", # 要执行的操作,比如暂停创建、调整Worker数量等
                    "data": new_mail # 邮件相关数据
                }
                self.report_queue.put(msg)

账号创建Worker示例

class AccountCreateWorker(multiprocessing.Process):
    def __init__(self, task_queue, report_queue):
        super().__init__()
        self.task_queue = task_queue # 任务下发队列,读取要执行的创建任务
        self.report_queue = report_queue # 上报队列,上报任务完成/错误状态
        self.running = True

    def run(self):
        while self.running:
            # 阻塞等待任务,超时1秒方便响应退出信号
            try:
                task = self.task_queue.get(timeout=1)
            except:
                continue
            if task.get("type") == "exit":
                self.running = False
                break
            # 执行账号创建逻辑
            account = self._create_account(task)
            # 上报完成状态
            self.report_queue.put({
                "type": "task_finish",
                "worker_id": self.pid,
                "account": account
            })

3. Handler类核心逻辑

class Handler:
    def __init__(self, worker_count):
        self.worker_count = worker_count
        # 创建通信队列
        self.report_queue = multiprocessing.Queue()
        self.task_queue = multiprocessing.Queue()
        # 启动IMAP监控进程
        self.imap_proc = IMAPProcess(self.report_queue)
        self.imap_proc.start()
        # 启动指定数量的Worker
        self.workers = []
        for _ in range(self.worker_count):
            worker = AccountCreateWorker(self.task_queue, self.report_queue)
            worker.start()
            self.workers.append(worker)
        # 先填充初始的账号创建任务
        local_account_count = self._get_local_account_count() # 替换为你自己的本地账号计数逻辑
        need_create = max(0, self.worker_count - local_account_count)
        for _ in range(need_create):
            self.task_queue.put({"type": "create_account"})

    def run(self):
        # 循环消费上报队列,处理所有子进程发来的消息
        while True:
            msg = self.report_queue.get()
            msg_type = msg.get("type")
            if msg_type == "new_mail":
                # 处理新邮件触发的操作,比如调整Worker数量、下发新任务等
                action = msg.get("action")
                if action == "add_worker":
                    new_worker = AccountCreateWorker(self.task_queue, self.report_queue)
                    new_worker.start()
                    self.workers.append(new_worker)
                elif action == "stop_create":
                    # 清空任务队列
                    while not self.task_queue.empty():
                        self.task_queue.get()
            elif msg_type == "task_finish":
                # 任务完成后如果账号还不够,继续下发创建任务
                local_account_count = self._get_local_account_count()
                if local_account_count < self.worker_count:
                    self.task_queue.put({"type": "create_account"})

4. 子进程间直接通信的实现

如果需要子进程之间不经过Handler直接通信,可以用multiprocessing.Manager创建共享的队列或者字典,所有进程都可以直接读写这个共享对象:

# 在Handler初始化时创建Manager
self.manager = multiprocessing.Manager()
# 创建子进程间共享的队列
self.shared_worker_queue = self.manager.Queue()
# 初始化子进程时把这个共享队列传进去,子进程之间就可以直接读写这个队列通信

注意事项

  • 所有需要跨进程传递的对象必须支持pickle序列化,套接字、进程实例这类不可序列化的对象不能写入队列
  • 建议所有消息都统一使用带type字段的字典格式,方便接收方快速识别消息类型做对应处理
  • 进程退出时建议先发送退出信号到队列,让子进程主动退出,避免资源泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 20:27:04