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

Python Azure Function消费计划下并行调用方法优化执行速度

解决方案

你的代码耗时过长存在两个核心问题,其中冗余资源加载是主要耗时来源,并行执行是次要优化手段,按以下优先级操作即可解决超时问题:

核心问题说明

  • 每次调用classify_mail都会重复执行模型加载、向量器初始化操作,这部分IO+反序列化操作占单次调用耗时的95%以上,属于完全冗余的开销
  • 现有代码中向量器调用fit_transform属于逻辑错误:推理阶段必须使用训练阶段拟合好的向量器执行transform,否则特征空间与训练时不一致,预测结果完全不可信
  • 分类预测属于CPU密集型任务,不要使用asyncio协程做并行,协程仅对IO密集型任务有效,CPU密集场景下协程不会产生提速效果,应使用多进程实现并行

第一步:预加载资源(优先级最高,可直接将总耗时从200秒级降到1秒级)

Azure Function的Python工作进程在冷启动完成后会持续存活,后续请求不会重置全局变量,因此可以在函数启动阶段一次性加载所有模型、向量器、配置项,无需每次请求、每次分类重复加载。

修改分类逻辑文件cfp.py

import os
import joblib
from concurrent.futures import ProcessPoolExecutor

# 全局缓存,函数实例生命周期内只加载一次
_MODEL_CACHE = {}
_VECTORIZER_CACHE = {}
# 按CPU核心数初始化全局进程池,避免每次请求重复创建进程的开销
_PROCESS_POOL = ProcessPoolExecutor(max_workers=min(os.cpu_count(), 4))
MODEL_ROOT = os.path.join(os.path.dirname(os.path.abspath(__file__)), "model")

def init_resources(selected_models, vectorizer_parameters):
    """冷启动时调用一次,预加载所有模型和训练好的向量器"""
    for model_type in selected_models.keys():
        for scenario in selected_models[model_type]['scenarios']:
            cache_key = f"{model_type}_{scenario}"
            # 替换为你实际的模型、向量器文件路径拼接规则
            model_path = os.path.join(MODEL_ROOT, model_type, scenario, "model.pkl")
            vec_path = os.path.join(MODEL_ROOT, model_type, scenario, "vectorizer.pkl")
            _MODEL_CACHE[cache_key] = joblib.load(model_path)
            _VECTORIZER_CACHE[cache_key] = joblib.load(vec_path)

def classify_mail(cache_key, X):
    """单场景分类逻辑,直接读取缓存资源执行预测"""
    model = _MODEL_CACHE[cache_key]
    vectorizer = _VECTORIZER_CACHE[cache_key]
    X_features = vectorizer.transform(X)
    return {"prediction": model.predict(X_features)[0]}

def predict(mail_cleaned, selected_models, thresholds):
    task_list = []
    cache_key_list = []
    # 批量提交预测任务到进程池
    for model_type in selected_models.keys():
        for scenario in selected_models[model_type]['scenarios']:
            cache_key = f"{model_type}_{scenario}"
            cache_key_list.append(cache_key)
            task_list.append(_PROCESS_POOL.submit(classify_mail, cache_key, [mail_cleaned]))
    # 收集所有并行执行的结果
    results = []
    for idx, task in enumerate(task_list):
        single_result = {"name": cache_key_list[idx]}
        single_result.update(task.result())
        results.append(single_result)
    
    prediction = {}
    # 保留你原有后续组装prediction的逻辑
    return prediction

修改入口文件__init__.py

将配置加载、资源初始化逻辑移到全局作用域,冷启动时只执行一次:

import os
# 提前配置线程数,避免numpy、sklearn默认多线程抢核降低并行效率
os.environ["OMP_NUM_THREADS"] = "1"
os.environ["OPENBLAS_NUM_THREADS"] = "1"
os.environ["MKL_NUM_THREADS"] = "1"

import json
import traceback
import azure.functions as func
import cfp

# 以下逻辑放到main函数外,函数启动时执行一次
# 加载classes、selected_models、thresholds、vectorizer_parameters等配置
cfp.init_resources(selected_models, vectorizer_parameters)

def main(req: func.HttpRequest, context: func.Context) -> func.HttpResponse:
    try:
        req_body = req.get_json()
    except ValueError:
        req_body = None

    if req_body:
        try:
            prediction = cfp.predict(
                req_body['text_cleaned'],
                selected_models,
                thresholds
            )
            return func.HttpResponse(
                json.dumps(prediction).encode('utf-8'),
                status_code=200,
                mimetype='application/json'
            )
        except Exception as e:
            return func.HttpResponse(
                json.dumps({'status': 'fehler', 'comment': str(e), 'stack_trace': traceback.format_exc()}),
                status_code=400,
                mimetype='application/json'
            )
    else:
        return func.HttpResponse(
            json.dumps({'status': 'fehler', 'comment': 'Die Eingabedaten wurden falsch angegeben', 'stack_trace': ''}),
            status_code=400,
            mimetype='application/json'
        )

第二步:配置调整

修改函数根目录下的host.json,匹配消费计划的资源限制:

{
  "version": "2.0",
  "extensionBundle": {
    "id": "Microsoft.Azure.Functions.ExtensionBundle",
    "version": "[3.*, 4.0.0)"
  },
  "functionTimeout": "00:05:00",
  "extensions": {
    "http": {
      "maxOutstandingRequests": 20,
      "maxConcurrentRequests": 10,
      "dynamicThrottlesEnabled": false
    }
  },
  "languageWorkers": {
    "python": {
      "maxProcessCount": 2
    }
  }
}
  • 消费计划默认单实例为2vCPU,maxProcessCount设为2即可,不要设置超过CPU核心数避免上下文切换开销
  • 消费计划最长超时时间可设置为10分钟,可根据实际预测耗时调整functionTimeout参数
  • 不要将入口函数定义为async def,CPU密集场景下同步入口+进程池的性能比异步入口高3-5倍,不会出现事件循环阻塞问题

内容的提问来源于stack exchange,提问作者sergioMoreno

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 01:03:33