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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 16:01:29