使用multiprocessing Pool的imap_unordered未完成全部任务问题排查
问题分析与解决方案
你的进程运行2-3小时后陷入“睡眠”,核心原因是单个API请求长时间无响应(卡住),导致对应的工作进程阻塞在API调用环节无法返回结果;其他进程完成任务后无新任务可处理,整体陷入停滞。此外代码存在几处潜在问题会加剧该情况:
- 缺乏API请求超时控制:若某个API无响应,工作进程会永久挂起,永远不会返回结果。
- 异常处理逻辑错误:
multiprocessing.Pool不会直接返回Exception实例,任务抛出的异常会在主进程迭代结果时触发,你的isinstance(result, Exception)判断无法捕获这类异常。 - 无进程重启机制:长时间运行的进程可能积累状态异常,导致意外挂起。
优化后的多进程代码
针对上述问题,调整代码如下,添加超时控制、完整异常捕获、进程重启机制及任务日志:
import logging from multiprocessing import Pool # 根据实际API调用库调整,此处假设用requests import requests # 配置进程级日志,便于排查卡住的任务 logging.basicConfig( level=logging.INFO, format="%(process)d - %(asctime)s - %(message)s" ) urls_to_query = [...] # 5000个URL列表 def call_api(url): logger = logging.getLogger(__name__) logger.info(f"启动任务: {url}") try: # 给API调用添加超时(示例设为300秒,根据实际响应时间调整) response = requests.get(url, timeout=300) # 替换为你的call_some_api实现 response.raise_for_status() # 捕获HTTP状态码异常 response_json = response.json() logger.info(f"完成任务: {url}") return {"success": True, "data": response_json, "url": url} except Exception as e: logger.error(f"任务失败 {url}: {str(e)}") return {"success": False, "error": str(e), "url": url} def write_to_file(result): # 追加写入JSON行,避免一次性加载大量数据 with open("output.json", "a", encoding="utf-8") as f: import json json.dump(result["data"], f) f.write("\n") def handle_failure(result): with open("errors.log", "a", encoding="utf-8") as f: f.write(f"{result['url']}: {result['error']}\n") if __name__ == "__main__": # maxtasksperchild设置为100,让每个进程处理100个任务后重启,避免状态异常 with Pool(processes=8, maxtasksperchild=100) as pool: total = len(urls_to_query) for i, result in enumerate(pool.imap_unordered(call_api, urls_to_query), 1): print(f"已处理 {i}/{total} - URL: {result['url']}") if result["success"]: write_to_file(result) else: handle_failure(result)
更优替代方案:IO密集型任务推荐用多线程/异步IO
你的任务属于IO密集型(大部分时间在等待API响应),多进程会浪费CPU资源,以下两种方案能大幅提升并发效率:
方案1:多线程(ThreadPoolExecutor)
线程开销远低于进程,可设置更多并发数(如32-64):
from concurrent.futures import ThreadPoolExecutor import logging import json import requests logging.basicConfig(level=logging.INFO, format="%(thread)d - %(asctime)s - %(message)s") urls_to_query = [...] def call_api(url): logger = logging.getLogger(__name__) logger.info(f"启动任务: {url}") try: response = requests.get(url, timeout=300) response.raise_for_status() response_json = response.json() logger.info(f"完成任务: {url}") return {"success": True, "data": response_json, "url": url} except Exception as e: logger.error(f"任务失败 {url}: {str(e)}") return {"success": False, "error": str(e), "url": url} def main(): total = len(urls_to_query) # 线程数设为CPU核心数的4-8倍,示例为32 with ThreadPoolExecutor(max_workers=32) as executor: for i, result in enumerate(executor.map(call_api, urls_to_query), 1): print(f"已处理 {i}/{total} - URL: {result['url']}") if result["success"]: with open("output.json", "a", encoding="utf-8") as f: json.dump(result["data"], f) f.write("\n") else: with open("errors.log", "a", encoding="utf-8") as f: f.write(f"{result['url']}: {result['error']}\n") if __name__ == "__main__": main()
方案2:异步IO(aiohttp)
异步IO能实现最高并发量,适合大规模API调用:
import aiohttp import asyncio import logging import json logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(message)s") urls_to_query = [...] async def call_api(session, url): logger = logging.getLogger(__name__) logger.info(f"启动任务: {url}") try: async with session.get(url, timeout=300) as response: response.raise_for_status() response_json = await response.json() logger.info(f"完成任务: {url}") return {"success": True, "data": response_json, "url": url} except Exception as e: logger.error(f"任务失败 {url}: {str(e)}") return {"success": False, "error": str(e), "url": url} async def main(): total = len(urls_to_query) # 控制并发数(根据API限流规则调整,示例为50) conn = aiohttp.TCPConnector(limit=50) async with aiohttp.ClientSession(connector=conn) as session: tasks = [call_api(session, url) for url in urls_to_query] for i, result_task in enumerate(asyncio.as_completed(tasks), 1): result = await result_task print(f"已处理 {i}/{total} - URL: {result['url']}") if result["success"]: with open("output.json", "a", encoding="utf-8") as f: json.dump(result["data"], f) f.write("\n") else: with open("errors.log", "a", encoding="utf-8") as f: f.write(f"{result['url']}: {result['error']}\n") if __name__ == "__main__": asyncio.run(main())
内容的提问来源于stack exchange,提问作者Karanvir Ghumaan
相关产品推荐
相关产品推荐

