基于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
相关产品推荐
相关产品推荐

