运行含Pyrogram的多进程代码时出现TypeError:无法pickle _queue.SimpleQueue对象
问题
运行以下Python代码时触发错误:TypeError: cannot pickle '_queue.SimpleQueue' object,需求是通过multiprocessing启动多个Pyrogram Client进程,每个进程并行解析用户信息后将结果发送到队列,最终汇总成完整的用户集合。
代码如下:
async def main(): print(message) queue = multiprocessing.Queue() apps = [] combs = batched(combinations, 200) chat_id = obtain_chats_id() workers = [] print(f'~~~Parsing channel {chat_id} ~~~') try: async for i in async_range(num_workers): async with Client(f'session_{i}', api_hash=api_hash, api_id=api_id) as app: apps.append(app) for app, comb in zip(apps, combs): worker = mp.Process(target=get_users_info, args=(app, queue, chat_id, comb)) workers.append(worker) worker.start() for worker in workers: worker.join() num_done = 0 while num_done < num_workers: msg = queue.get() if msg == "DONE": num_done += 1 else: final_set.update(msg) finally: save_excel()
错误原因
- Pyrogram Client无法跨进程序列化:你在主进程中创建Client实例后,直接将其作为参数传递给子进程。但Pyrogram Client内部包含网络连接、内部队列等不可被
pickle序列化的资源,而multiprocessing启动子进程时必须对传入参数做序列化操作,这直接触发了报错。 - 隐含的队列序列化问题:即使你显式使用了
multiprocessing.Queue,Client内部自带的_queue.SimpleQueue对象无法被序列化,跟着Client一起传递时暴露了这个问题。
解决方案
方案1:子进程内独立创建Client
不在主进程初始化Client,而是将创建Client所需的参数(会话名、api_id、api_hash)传给子进程,让子进程自行初始化Client实例:
async def main(): print(message) queue = multiprocessing.Queue() combs = batched(combinations, 200) chat_id = obtain_chats_id() workers = [] print(f'~~~Parsing channel {chat_id} ~~~') try: # 直接生成子进程需要的Client创建参数 for i, comb in enumerate(combs): worker = mp.Process( target=get_users_info, args=(f'session_{i}', api_id, api_hash, queue, chat_id, comb) ) workers.append(worker) worker.start() for worker in workers: worker.join() # 汇总结果逻辑不变 num_done = 0 while num_done < len(workers): msg = queue.get() if msg == "DONE": num_done += 1 else: final_set.update(msg) finally: save_excel() # 修改get_users_info,在子进程内部创建并使用Client def get_users_info(session_name, api_id, api_hash, queue, chat_id, comb): from pyrogram import Client import asyncio async def run_parse(): async with Client(session_name, api_id=api_id, api_hash=api_hash) as app: # 这里编写你的用户信息解析逻辑 user_set = set() # ... 解析代码 ... queue.put(user_set) queue.put("DONE") asyncio.run(run_parse())
方案2:改用多线程(IO密集型场景适用)
如果你的解析工作以网络IO为主,可以用multiprocessing.dummy(即多线程)替代多进程,线程共享进程内存空间,无需序列化Client实例:
from multiprocessing.dummy import Pool as ThreadPool from functools import partial async def main(): print(message) queue = multiprocessing.Queue() apps = [] combs = batched(combinations, 200) chat_id = obtain_chats_id() print(f'~~~Parsing channel {chat_id} ~~~') try: async for i in async_range(num_workers): async with Client(f'session_{i}', api_hash=api_hash, api_id=api_id) as app: apps.append(app) # 创建线程池,直接传递Client实例 pool = ThreadPool(num_workers) worker_func = partial(get_users_info, queue=queue, chat_id=chat_id) pool.starmap(worker_func, zip(apps, combs)) pool.close() pool.join() # 汇总结果逻辑不变 num_done = 0 while num_done < num_workers: msg = queue.get() if msg == "DONE": num_done += 1 else: final_set.update(msg) finally: save_excel()
关键注意事项
- 禁止在进程间传递Pyrogram Client实例,每个进程必须拥有独立的Client。
- 子进程内的异步代码需要用
asyncio.run()启动,因为子进程没有默认的事件循环。 - 多线程场景下,建议每个线程使用独立的Client实例,避免并发操作冲突。
内容的提问来源于stack exchange,提问作者Emelian Pugachev
相关产品推荐
相关产品推荐

