انتقل إلى المحتوى الرئيسي

الوحدة 6 — الطلبات غير المتزامنة ومهامّ الخلفيّة

async كلمة سحريّة كثيرًا ما تُستعمل بلا فهم عواقبها. في خدمة تنبّؤ، إضافة async أمام كلّ دالّة تعتقد أنّها ستُسرّع الأمور، وقد تُبطئها أضعافًا. هذه الوحدة تفرّق بين الحالة التي تستفيد فيها فعلًا من async والحالة التي يجب فيها تجنّبها، ثمّ تُقدّم مهامّ الخلفيّة للأعمال الطويلة (تسجيل ملفّ كامل) التي لا يمكن للعميل انتظارها.

async في جملتَين

FastAPI يعتمد على حلقة أحداث واحدة (asyncio event loop) لكلّ عامل. async def يُتيح تعليق الدالّة عند عمليّة I/O (استعلام قاعدة بيانات، استدعاء HTTP لخدمة أخرى، قراءة قرص) لتُعالج الحلقة طلبًا آخر. بدون async، الدالّة تشغل الخيط حتّى تنتهي.

القاعدة: async يُفيد عندما تكون الدالّة في انتظار I/O بلا حساب، ولا يُفيد (بل يضرّ) عندما تكون في حساب متواصل يستهلك CPU (كتنبّؤ نموذج).

المسار المتزامن العاديّ

FastAPI يُميّز تلقائيًّا بين مسار def عاديّ ومسار async def. المسار العاديّ يُنفَّذ في مجمع خيوط (thread pool) مُنفصل عن حلقة الأحداث، فلا يعطّلها:

@app.post("/predict", response_model=ChurnResponse)
def predire(req: ChurnRequest) -> ChurnResponse:
# يُنفَّذ في خيط منفصل ; الحلقة تظلّ تعالج طلبات أخرى
df = pd.DataFrame([req.model_dump()])
proba = float(app.state.model.predict_proba(df)[0][1])
...

هذا هو الاختيار الافتراضيّ الصحيح لكلّ مسار تنبّؤ. النموذج يستهلك CPU (وربّما GIL)، والخيط المنفصل يمنعه من تجميد الخدمة كاملةً.

الفخّ: async def مع نموذج حاجب

الشيفرة التالية تبدو أنيقة لكنّها تدمّر الخدمة:

@app.post("/predict", response_model=ChurnResponse)
async def predire(req: ChurnRequest) -> ChurnResponse:
# كارثة : يعطّل الحلقة كاملةً
df = pd.DataFrame([req.model_dump()])
proba = float(app.state.model.predict_proba(df)[0][1])
...

المشكلة: predict_proba استدعاء متزامن حاجب يشغّل CPU لعشرات الميلّي ثواني. داخل async def يعمل مباشرة في حلقة الأحداث الوحيدة للعامل. ما دام يعمل، كلّ طلبات العميل الأخرى معلّقة. تحت حمل، الإنتاجيّة تنخفض ×10 مقارنةً بالمسار المتزامن العاديّ الذي يعمل في مجمع خيوط.

القاعدة القطعيّة: لا تُضِف async أمام مسار يستدعي نموذجًا حاجبًا. إمّا أن تُبقيه def عاديًّا، أو تُغلّف الاستدعاء بـasyncio.to_thread:

import asyncio


@app.post("/predict", response_model=ChurnResponse)
async def predire(req: ChurnRequest) -> ChurnResponse:
df = pd.DataFrame([req.model_dump()])
proba = await asyncio.to_thread(
lambda: float(app.state.model.predict_proba(df)[0][1])
)
...

هذا يُعطي نفس فائدة الخيط المنفصل، ويسمح باستعمال async لأشياء أخرى في نفس الدالّة (استعلامات I/O مثلًا).

متى async يفيد فعلًا

مسار يستعلم قاعدة بيانات ثمّ يستدعي خدمة خارجيّة، بلا حساب مركّز:

import httpx


@app.get("/predict/enriched/{cid}")
async def predire_enrichi(cid: str) -> dict:
async with httpx.AsyncClient() as client:
req = await client.get(f"https://crm.example/customer/{cid}", timeout=2.0)
donnees = req.json()

df = pd.DataFrame([donnees])
proba = await asyncio.to_thread(
lambda: float(app.state.model.predict_proba(df)[0][1])
)
return {"customer_id": cid, "churn_probability": proba}

الاستعلام الخارجيّ ينتظر شبكةً لعدّة ميلّي ثوانٍ. async يسمح للعامل بمعالجة طلبات أخرى أثناء الانتظار. تحت حمل، الإنتاجيّة قد تتضاعف أربع مرّات.

مهامّ الخلفيّة: التسجيل الفوريّ لملفّ كامل

سيناريو حقيقيّ: العميل يرفع ملفّ CSV بخمسين ألف مشترك ويطلب تنبّؤًا للكلّ. المعالجة تستغرق دقيقتَين. لا يُعقل أن ينتظر HTTP client دقيقتَين مفتوحًا. الحلّ: مهمّة خلفيّة تُعيد فورًا معرّف تتبّع، والعميل يستطلع الحالة لاحقًا.

FastAPI يوفّر BackgroundTasks للأعمال القصيرة (ثوانٍ)، ونحتاج طابور خارجيّ (Celery، RQ، Dramatiq) للأعمال الطويلة (دقائق).

BackgroundTasks للتسجيل السريع

from fastapi import BackgroundTasks
import shutil
import uuid


def scorer_fichier(chemin: str, job_id: str, modele) -> None:
df = pd.read_csv(chemin)
df["churn_probability"] = modele.predict_proba(df)[:, 1]
df.to_csv(f"jobs/{job_id}.csv", index=False)


@app.post("/predict/file")
async def predire_fichier(
fichier: UploadFile,
background: BackgroundTasks,
) -> dict:
job_id = str(uuid.uuid4())
dest = f"jobs/upload-{job_id}.csv"
with open(dest, "wb") as f:
shutil.copyfileobj(fichier.file, f)

background.add_task(scorer_fichier, dest, job_id, app.state.model)
return {"job_id": job_id, "status": "queued"}

المهمّة تُنفَّذ بعد إرسال الاستجابة إلى العميل. العميل يحصل على job_id فورًا (202 Accepted لو أردنا الدقّة)، ويستطلع مسارًا GET /jobs/{job_id} لاحقًا للحصول على النتيجة.

قيد BackgroundTasks: المهمّة تعمل في نفس عمليّة العامل. إن أُعيد تشغيل الخدمة (نشر، تحديث)، المهمّة تُفقَد. وإن استغرقت دقائق، تشغل موارد العامل. للأعمال الطويلة أو الحرجة، ننتقل إلى طابور خارجيّ.

طابور خارجيّ: Celery أو RQ

للتسجيل الليليّ لملفّ بمليون سجلّ، البنية الأنسب:

# tasks.py
from celery import Celery

celery_app = Celery("churn", broker="redis://redis:6379/0")


@celery_app.task(name="score_batch")
def score_batch(job_id: str, chemin: str) -> str:
modele = joblib.load("churn_pipeline.joblib")
df = pd.read_csv(chemin)
df["proba"] = modele.predict_proba(df)[:, 1]
dest = f"jobs/{job_id}-out.csv"
df.to_csv(dest, index=False)
return dest


# main.py (FastAPI)
from tasks import score_batch


@app.post("/predict/file/async")
async def scoring_asynchrone(fichier: UploadFile) -> dict:
job_id = str(uuid.uuid4())
dest = f"jobs/upload-{job_id}.csv"
with open(dest, "wb") as f:
shutil.copyfileobj(fichier.file, f)

task = score_batch.delay(job_id, dest)
return {"job_id": job_id, "celery_id": task.id, "status": "queued"}

المزايا: عمّال Celery عمليّات مستقلّة، إعادة تشغيل FastAPI لا تُفقد المهمّة، عدد العمّال يتوسّع مستقلًّا عن عدد عمّال HTTP، وسجلّ المهام محفوظ في Redis.

قاعدة قرار مختصرة

الحالةالاختيار
تنبّؤ فرديّ سريعdef عاديّ (خيط منفصل تلقائيًّا)
تنبّؤ + استعلام I/Oasync def + asyncio.to_thread للنموذج
تسجيل ملفّ خفيف (ثوانٍ)BackgroundTasks
تسجيل ملفّ ثقيل (دقائق+)Celery/RQ + Redis
مسار مصحوب بـasync def بلا awaitخطأ تصميم، احذف async

الخلاصة

  • async def مع استدعاء نموذج حاجب يعطّل حلقة الأحداث ويُدمّر الإنتاجيّة؛ إمّا def عاديّ أو asyncio.to_thread.
  • المسار def العاديّ يُنفَّذ في مجمع خيوط، وهذا اختيار افتراضيّ آمن لكلّ مسار تنبّؤ.
  • BackgroundTasks مناسب للأعمال القصيرة داخل نفس العمليّة؛ للأعمال الطويلة أو الحرجة، طابور خارجيّ (Celery + Redis).
  • async يُفيد فعلًا عند انتظار I/O (قواعد بيانات، استدعاءات HTTP خارجيّة) لا عند حساب CPU مركّز.