المهام الخلفية وCelery — طوابير المهام غير المتزامنة
BackgroundTasks مثل نظام أرقام الطابور في المطعم — تحصل على رقم (استجابة API) فور الطلب، ويُبلَّغك عندما يكون طعامك جاهزًا؛ أما Celery فكمطبخ مركزي — عدة طهاة يعدّون طلبات مختلفة في وقت واحد، ويمكن تتبع حالة كل طلب، مع إعادة المحاولة عند الفشل.
1. ما ستتعلمه
BackgroundTasksالمدمج في FastAPI: حل سريع للسيناريوهات البسيطة- معمارية Celery: Worker / Broker (Redis) / Backend / مراقبة Flower
- دمج Celery مع FastAPI: تعريفات المهام، المشغلات، ونقاط نهاية استعلام الحالة
- إعادة محاولة المهام ومعالجة الأخطاء:
@app.task(retry=3, acks_late=True) - سيناريو Alice: جمع الأسعار الدفعة — تُعيد API معرّف مهمة فورًا، ويعالج Celery Worker بشكل غير متزامن جمع أسعار ملايين المنتجات
2. القصة الحقيقية لـ Alice
(1) المشكلة: المهام المستهلكة للوقت تحظر استجابات API
يحتاج PriceTracker الخاص بـ Alice إلى جمع أسعار ملايين المنتجات دفعة واحدة، وتستغرق العملية الواحدة 30 دقيقة. إذا تمت المعالجة بشكل متزامن، ستنتهي مهلة طلبات API، وسيضطر واجهة Bob الأمامية للانتظار 30 دقيقة لتلقي استجابة — مما يؤدي لتجربة مستخدم سيئة. في هذه الأثناء، يمكن لـ BackgroundTasks في FastAPI العمل فقط ضمن العملية الحالية، وإعادة تشغيل worker ستؤدي لفقدان المهمة.
(2) حل طوابير Celery الموزعة
يضع Celery المهام المستهلكة للوقت في طابور رسائل (Redis Broker)، حيث تستهلك عمليات Worker المهام بشكل غير متزامن، وتعيد API معرّف المهمة فورًا. تُعاد المهام تلقائيًا عند فشلها، ويمكن توسيع Workers أفقيًا، وإعادة تشغيل العمليات لا تؤثر على المهام في الطابور.
from celery import Celery
celery_app = Celery("pricetracker", broker="redis://localhost:6379/0")
@celery_app.task(bind=True, max_retries=3)
def scrape_prices(self, product_ids: list[int]):
# جمع أسعار غير متزامن - يعمل في Celery Worker
...
(3) العائد
انتقل جمع مليون منتج من "حظر API لـ 30 دقيقة" إلى "إعادة معرّف مهمة في ثانية واحدة + معالجة خلفية". يمكن توسيع عدد workers من 1 إلى 10، مما يقلل وقت الجمع من 30 دقيقة إلى 3 دقائق. تُعاد المهام تلقائيًا ثلاث مرات عند الفشل، مما يرفع معدل النجاح من 95% إلى 99.9%.
3. الحلول الخفيفة للمهام الخلفية
(1) السيناريوهات القابلة للتطبيق
(1) ▶ مثال: إرسال الإشعارات باستخدام BackgroundTasks
from fastapi import FastAPI, BackgroundTasks
from pydantic import BaseModel
app = FastAPI()
class PriceAlertRequest(BaseModel):
product_id: int
target_price: float
email: str
def send_price_alert_email(email: str, product_id: int, price: float):
# محاكاة إرسال بريد إلكتروني (لا تستخدم await هنا)
print(f"Sending alert to {email}: Product {product_id} hit ${price}")
@app.post("/alerts")
async def create_alert(alert: PriceAlertRequest, bg: BackgroundTasks):
# إضافة مهمة تعمل بعد إرسال الاستجابة
bg.add_task(send_price_alert_email, alert.email, alert.product_id, alert.target_price)
return {"message": "Alert created", "product_id": alert.product_id}
الناتج:
# تم تعريف الدالة بنجاح
(2) شجرة القرار: BackgroundTasks مقابل Celery
flowchart TD
Start{هل تحتاج مهمة خلفية؟} --> Time{تستغرق > 1 دقيقة؟}
Time -->|لا| Simple[استخدم BackgroundTasks]
Time -->|نعم| Retry{هل تحتاج إعادة محاولة/مرونة؟}
Retry -->|لا| Simple
Retry -->|نعم| Scale{هل تحتاج توسع أفقي؟}
Scale -->|لا| Simple
Scale -->|نعم| Celery[استخدم Celery]
Simple -->|المميزات| P1[بسيط، بدون بنية تحتية]
Simple -->|العيوب| C1[بدون إعادة محاولة، بدون توسع، تُفقد عند إعادة التشغيل]
Celery -->|المميزات| P2[إعادة محاولة، توسع، استمرار، مراقبة]
Celery -->|العيوب| C2[بنية تحتية Redis + Worker]
| البُعد | المهام الخلفية | Celery |
|---|---|---|
| التعقيد | بدون إعداد | يتطلب Redis + Worker |
| الاستمرار | ذاكرة العملية | استمرار Redis |
| إعادة المحاولة | لا يوجد | آلية إعادة محاولة مدمجة |
| التوسع | عملية واحدة | توسع أفقي قائم على Worker |
| المراقبة | لا يوجد | لوحة Flower |
| حالة الاستخدام | مهام خفيفة < 1 ثانية | مهام مستهلكة للوقت > 1 دقيقة |
4. معمارية Celery
(1) نظرة عامة على المعمارية
flowchart TD
API[تطبيق FastAPI] -->|إدراج مهمة| Broker[(Redis Broker)]
Broker -->|استهلاك مهمة| Worker1[Celery Worker 1]
Broker -->|استهلاك مهمة| Worker2[Celery Worker 2]
Broker -->|استهلاك مهمة| WorkerN[Celery Worker N]
Worker1 -->|تخزين النتيجة| Backend[(Redis Backend)]
Worker2 -->|تخزين النتيجة| Backend
WorkerN -->|تخزين النتيجة| Backend
API -->|استعلام الحالة| Backend
Flower[مراقب Flower] -->|مراقبة| Broker
Flower -->|مراقبة| Backend
| المكون | الوظيفة | التوصية |
|---|---|---|
| Broker | طابور رسائل المهام | Redis |
| Backend | تخزين النتائج | Redis |
| Worker | عملية تنفيذ المهام | celery -A app worker |
| Flower | لوحة المراقبة | celery -A app flower |
(1) ▶ مثال: إعدادات Celery
# app/core/celery_app.py
from celery import Celery
celery_app = Celery(
"pricetracker",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
task_track_started=True,
task_acks_late=True, # إقرار بعد التنفيذ، وليس قبله
worker_prefetch_multiplier=4,
result_expires=3600, # تنتهي صلاحية النتائج بعد ساعة واحدة
)
الناتج:
# تم التنفيذ بنجاح
5. تعريف مهام Celery وتشغيلها
(1) تعريف المهمة
(1) ▶ مثال: مهمة جمع الأسعار
# app/tasks/price_scraping.py
from app.core.celery_app import celery_app
import asyncio
from sqlalchemy import select
@celery_app.task(bind=True, max_retries=3, default_retry_delay=60)
def scrape_product_prices(self, product_ids: list[int]):
"""جمع الأسعار الحالية للمنتجات المحددة."""
try:
for pid in product_ids:
# محاكاة الجمع (في الإنتاج: طلبات HTTP لمصادر الأسعار)
price = fetch_price_from_source(pid)
# التخزين في قاعدة البيانات
save_price_record(pid, price)
return {"scraped": len(product_ids), "status": "success"}
except Exception as exc:
# إعادة المحاولة مع تراجع أسي
raise self.retry(exc=exc, countdown=60 * (2 ** self.request.retries))
@celery_app.task(bind=True)
def bulk_scrape_all(self, total_products: int = 1000000, batch_size: int = 1000):
"""جمع جميع المنتجات على دفعات - نمط chord."""
batches = [
list(range(i, min(i + batch_size, total_products)))
for i in range(0, total_products, batch_size)
]
# توزيع على مهام جمع فردية
for batch in batches:
scrape_product_prices.delay(batch)
return {"total_batches": len(batches), "status": "started"}
الناتج:
# تم تعريف الدالة بنجاح
(2) انتقالات حالة المهمة
stateDiagram-v2
[*] --> PENDING: تم إنشاء المهمة
PENDING --> STARTED: التقط Worker المهمة
STARTED --> PROGRESS: قيد التشغيل (اختياري)
PROGRESS --> SUCCESS: اكتملت
PROGRESS --> FAILURE: حدث خطأ
STARTED --> FAILURE: حدث خطأ
FAILURE --> RETRY: لم يصل max_retries
RETRY --> PENDING: أُعيد إدراجها
FAILURE --> [*]: تجاوز max_retries
SUCCESS --> [*]
(2) ▶ مثال: نقطة نهاية FastAPI تشغل مهمة Celery
from fastapi import FastAPI, Depends
from app.core.celery_app import celery_app
from app.tasks.price_scraping import scrape_product_prices, bulk_scrape_all
app = FastAPI()
@app.post("/api/v1/scrape/prices")
async def trigger_scrape(
product_ids: list[int],
user=Depends(require_subscription("pro")),
):
# تشغيل مهمة Celery - تُعيد معرّف المهمة فورًا
task = scrape_product_prices.delay(product_ids)
return {"task_id": task.id, "status": "pending"}
@app.post("/api/v1/scrape/bulk")
async def trigger_bulk_scrape(
user=Depends(require_subscription("enterprise")),
):
task = bulk_scrape_all.delay(total_products=1000000)
return {"task_id": task.id, "status": "pending"}
الناتج:
# تم تعريف الدالة بنجاح
(3) ▶ مثال: نقطة نهاية استعلام حالة المهمة
from celery.result import AsyncResult
@app.get("/api/v1/tasks/{task_id}")
async def get_task_status(task_id: str):
result = AsyncResult(task_id, app=celery_app)
response = {
"task_id": task_id,
"status": result.status,
}
if result.ready():
if result.successful():
response["result"] = result.result
else:
response["error"] = str(result.result)
elif result.state == "PROGRESS":
response["progress"] = result.info
return response
الناتج:
# تم تعريف الدالة بنجاح
6. تشغيل ومراقبة Celery Workers
(1) ▶ مثال: بدء Worker وFlower
# بدء Celery Worker
celery -A app.core.celery_app worker --loglevel=info --concurrency=4
# بدء لوحة مراقبة Flower
celery -A app.core.celery_app flower --port=5555
# زيارة http://localhost:5555 للوحة المراقبة
الناتج:
# تم تنفيذ الأمر بنجاح
(2) ▶ مثال: تقرير تقدم المهمة
from celery import current_task
@celery_app.task(bind=True)
def scrape_with_progress(self, product_ids: list[int]):
total = len(product_ids)
for i, pid in enumerate(product_ids):
# معالجة كل منتج
price = fetch_price_from_source(pid)
save_price_record(pid, price)
# الإبلاغ عن التقدم
self.update_state(
state="PROGRESS",
meta={"current": i + 1, "total": total, "percent": (i + 1) / total * 100},
)
return {"scraped": total}
الناتج:
# تم تعريف الدالة بنجاح
❓ أسئلة شائعة
task_acks_late=True؟acks_late=True يجعل Worker يقرّ بالمهمة فقط بعد اكتمال التنفيذ؛ إذا تعطل، ستُنفَّذ المهمة بواسطة worker آخر.async/await في مهام Celery؟asyncio.run(). أو استخدم نماذج التزامن eventlet أو gevent في Worker.📖 ملخص
- BackgroundTasks مناسبة للمهام الخفيفة التي تستغرق أقل من ثانية (إرسال رسائل بريد إلكتروني، كتابة سجلات)؛ لا تتطلب إعدادًا لكنها لا تدعم إعادة المحاولة أو الاستمرار.
- Celery طابور مهام بمستوى الإنتاج: Broker (طابور Redis) + Worker (تنفيذ) + Backend (تخزين النتائج)
@app.task(bind=True, max_retries=3)عرّف مهام قابلة لإعادة المحاولة؛self.retry()شغّل إعادة المحاولة- FastAPI يشغل المهمة عبر
.delay()ويستعلم عن الحالة عبرAsyncResult؛ تُعيد API معرّف المهمة فورًا - Flower يوفر لوحة مراقبة بصرية تتيح لك عرض حالة workers وتقدم المهام ومعدلات النجاح في الوقت الفعلي
📝 تمارين
- مسألة أساسية (الصعوبة ⭐): استخدم
BackgroundTasksلتنفيذ نقطة نهاية تُرسل بريد إلكتروني إشعار في الخلفية (محاكاة إخراج السجل) بعد إنشاء منتج، وتُعيد API نتيجة الإنشاء فورًا. تلميح:bg: BackgroundTasks+bg.add_task(fn, args) - تمرين متقدم (الصعوبة ⭐⭐): اضبط Celery (Redis Broker + Backend)، عرّف مهمة
scrape_prices(max_retries=3)، اجعل نقطة نهاية FastAPI تشغل المهمة وتُعيد task_id، واجعل نقطة نهاية أخرى تستعلم عن حالة المهمة. تلميح:celery_app.delay()+AsyncResult(task_id) - تحدٍ (الصعوبة: ⭐⭐⭐): نفّذ مهمة جمع دفعات مع تقرير تقدم —
scrape_with_progressيُبلغ بنسبة التقدم إلىself.update_state(state="PROGRESS")، الواجهة الأمامية تستقصي/tasks/{task_id}لعرض شريط تقدم، ومستخدمو Enterprise يمكنهم تشغيل جمع بمستوى المليون. تلميح: استخدمself.update_state()+result.infoلاسترجاع التقدم
---|



