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
相关产品推荐
相关产品推荐

