Python多进程百万接口请求:语法疑问与性能优化咨询
Python多进程处理百万WebService请求问题解答
疑问解答
1. 两种进程池初始化语法的区别
- 直接初始化
Pool:需手动调用pool.close()(禁止提交新任务)和pool.join()(等待所有子进程完成),否则可能出现资源泄漏或子进程未结束主程序就退出的问题。 with Pool(...) as pool:依托Python上下文管理器特性,代码块结束时自动执行close()和join(),无需手动维护,更简洁安全,能避免遗忘操作导致的异常。
2. 能否逐行读取ID并立即提交任务
完全可以,无需将百万ID全部加载到内存再提交,这样能大幅节省内存占用。推荐使用pool.imap或pool.imap_unordered(不关心结果顺序时),以迭代器方式逐行提交任务:
# 包装函数,适配imap的单参数要求 def starmap_wrapper(args): return task(*args) with Pool(processes=4) as pool: with open(list_of_ids, 'r') as infile: # 用生成器逐行生成参数,无需存全量数据 task_args = ((line.strip(), queue) for line in infile) results = pool.imap(starmap_wrapper, task_args) # 遍历惰性生成的结果,触发任务执行并等待完成 for _ in results: pass
3. 避免重复传递大的requestBodyTemplate
推荐通过进程池初始化函数,在子进程启动时一次性传入模板,避免每个任务重复传递大对象:
# 定义全局变量存储模板,供子进程任务调用 global_request_body = None def init_worker(template): global global_request_body global_request_body = template # 初始化进程池时传入初始化函数和模板参数 with Pool(processes=4, initializer=init_worker, initargs=(requestBodyTemplate,)) as pool: # 任务函数直接使用全局变量global_request_body,无需再传该参数 def task(i, queue): # 基于模板生成请求体 request_body = global_request_body.replace("{ID}", i.strip()) ...
4. 确保所有子任务完成后主程序退出
遵循以下步骤即可:
- 所有任务提交完成后,调用
pool.close(),禁止向进程池新增任务; - 调用
pool.join(),阻塞主进程直到所有子进程任务执行完毕。
如果使用with语句初始化进程池,会自动执行上述两步;如果手动初始化Pool,必须按顺序调用这两个方法。若使用apply_async逐个提交任务,需保存所有AsyncResult对象,遍历调用res.get()等待结果:
pool = Pool(processes=4) results = [] with open(list_of_ids, 'r') as infile: for line in infile: res = pool.apply_async(task, args=(line.strip(), queue)) results.append(res) # 等待所有任务完成 for res in results: res.get() pool.close() pool.join()
5. 为任务设置超时时间
推荐使用concurrent.futures.ProcessPoolExecutor,它支持更灵活的超时控制,超时会抛出TimeoutError:
from concurrent.futures import ProcessPoolExecutor, TimeoutError with ProcessPoolExecutor(max_workers=4) as executor: # 提交所有任务,保存future与对应ID的映射 futures = {} with open(list_of_ids, 'r') as infile: for line in infile: i = line.strip() futures[executor.submit(task, i, queue)] = i # 遍历处理结果 for future in futures: try: # 设置30秒超时 result = future.result(timeout=30) except TimeoutError: print(f"ID {futures[future]} 任务超时") except Exception as e: print(f"ID {futures[future]} 任务出错: {str(e)}")
若使用multiprocessing.Pool,可在apply_async的get()方法中设置超时,超时后可调用pool.terminate()强制终止所有进程(会中断所有运行中任务,需谨慎使用)。
其他优化建议
- 合理设置进程数:WebService请求属于IO密集型任务,进程数建议设为CPU核心数的2-4倍,过多进程会导致网络连接过载,反而降低效率;
- 复用HTTP连接:在子进程初始化时创建
requests.Session,复用HTTP连接,避免每个任务重复建立连接的开销; - 状态记录优化:用
filelock库实现线程安全的状态文件写入,或改用SQLite数据库记录任务状态,避免队列+监听进程的通信瓶颈; - 添加异常重试:对网络错误、服务器5xx错误等场景添加重试逻辑,可使用
tenacity库简化重试代码; - 流式处理响应:如果响应内容较大,直接流式写入文件,避免全部加载到内存;
- 进度监控:用
tqdm库显示任务进度,方便实时了解执行情况:
from tqdm import tqdm with open(list_of_ids, 'r') as infile: ids = [line.strip() for line in infile] with Pool(processes=4) as pool: for _ in tqdm(pool.imap(task, ((i, queue) for i in ids)), total=len(ids)): pass
现有代码修正示例
针对你提供的代码中的问题(如queue未定义就使用、冗余sleep等),修正后的核心代码如下:
from multiprocessing import Pool, Manager, Process, current_process, set_start_method import requests global_request_body = None session = None def init_worker(template): global global_request_body, session global_request_body = template session = requests.Session() def listener(queue): with open('task_status.txt', 'a') as f: while True: msg = queue.get() if msg is None: break f.write(f"{msg}\n") f.flush() def task(i, queue): process = current_process() try: request_body = global_request_body.replace("{ID}", i) response = session.post("https://your-service-url.com", json=request_body, timeout=30) response.raise_for_status() with open(f"./responses/{i}.json", 'w') as f: f.write(response.text) queue.put(f"SUCCESS - {i} - {process.name}") except Exception as e: queue.put(f"FAILED - {i} - {process.name} - {str(e)}") def getRequestBodyTemplateJSON(): return '{"id": "{ID}", "data": "sample"}' def main(): set_start_method('spawn') list_of_ids = 'ids.txt' requestBodyTemplate = getRequestBodyTemplateJSON() with Manager() as manager: queue = manager.Queue() listener_proc = Process(target=listener, args=(queue,)) listener_proc.start() with Pool(processes=4, initializer=init_worker, initargs=(requestBodyTemplate,)) as pool: with open(list_of_ids, 'r') as infile: results = pool.imap(task, ((line.strip(), queue) for line in infile)) for _ in results: pass queue.put(None) listener_proc.join() if __name__ == "__main__": main()
内容的提问来源于stack exchange,提问作者coredump
相关产品推荐
相关产品推荐

