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

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. 确保所有子任务完成后主程序退出

遵循以下步骤即可:

  1. 所有任务提交完成后,调用pool.close(),禁止向进程池新增任务;
  2. 调用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()强制终止所有进程(会中断所有运行中任务,需谨慎使用)。

其他优化建议

  1. 合理设置进程数:WebService请求属于IO密集型任务,进程数建议设为CPU核心数的2-4倍,过多进程会导致网络连接过载,反而降低效率;
  2. 复用HTTP连接:在子进程初始化时创建requests.Session,复用HTTP连接,避免每个任务重复建立连接的开销;
  3. 状态记录优化:用filelock库实现线程安全的状态文件写入,或改用SQLite数据库记录任务状态,避免队列+监听进程的通信瓶颈;
  4. 添加异常重试:对网络错误、服务器5xx错误等场景添加重试逻辑,可使用tenacity库简化重试代码;
  5. 流式处理响应:如果响应内容较大,直接流式写入文件,避免全部加载到内存;
  6. 进度监控:用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 07:18:12