پردازش غیرهم‌زمان API هوش مصنوعی؛ آموزش Job Queue، Worker، Polling و Webhook

در این آموزش یک معماری غیرهم‌زمان برای API هوش مصنوعی می‌سازیم که درخواست‌ها را با FastAPI دریافت می‌کند، در Redis و Celery پردازش می‌کند و نتیجه را از طریق Polling یا Webhook تحویل می‌دهد.

Share
پردازش غیرهم‌زمان API هوش مصنوعی؛ آموزش Job Queue، Worker، Polling و Webhook

مقدمه

در اولین نسخه یک محصول مبتنی بر هوش مصنوعی، معمولاً درخواست کاربر مستقیماً از Backend به مدل ارسال می‌شود و برنامه تا دریافت پاسخ منتظر می‌ماند:

کاربر → Backend → API هوش مصنوعی → مدل → پاسخ

این معماری برای درخواست‌های کوتاه متنی قابل‌قبول است، اما در سناریوهای زیر به‌سرعت مشکل‌ساز می‌شود:

  • تولید تصویر با کیفیت بالا
  • ساخت ویدئو با هوش مصنوعی
  • تبدیل متن به گفتار طولانی
  • تحلیل فایل‌های بزرگ
  • پردازش دسته‌ای هزاران رکورد
  • تولید گزارش‌های مفصل
  • اجرای گردش‌کار چندمرحله‌ای
  • پردازش اسناد با OCR و مدل زبانی
  • اجرای چند فراخوانی متوالی API
  • وظایف ایجنتی که به چند ابزار متصل می‌شوند

در چنین شرایطی ممکن است هر درخواست از چند ثانیه تا چند دقیقه زمان ببرد. نگه‌داشتن اتصال HTTP در تمام این مدت باعث افزایش Timeout، مصرف منابع سرور، تجربه کاربری ضعیف و دشوارشدن مقیاس‌پذیری می‌شود.

راه‌حل استاندارد، جداکردن «ثبت درخواست» از «اجرای پردازش» است. در این معماری، Backend درخواست را دریافت می‌کند، یک Job می‌سازد، آن را در صف قرار می‌دهد و بلافاصله شناسه Job را به کاربر برمی‌گرداند. سپس Workerها وظیفه را در پس‌زمینه اجرا می‌کنند.

در این مقاله یک سیستم عملی با اجزای زیر می‌سازیم:

  • FastAPI برای ارائه REST API
  • Celery برای مدیریت Task و Worker
  • Redis برای Message Broker و Result Backend
  • Polling برای دریافت وضعیت Job
  • Webhook برای اعلام خودکار نتیجه
  • API درواره برای اجرای مدل هوش مصنوعی
  • Idempotency Key برای جلوگیری از ثبت ناخواسته درخواست تکراری
  • Retry و Exponential Backoff برای خطاهای موقت
  • Docker Compose برای اجرای ساده سرویس‌ها

پردازش هم‌زمان و غیرهم‌زمان چه تفاوتی دارند؟

پردازش هم‌زمان یا Synchronous

در روش هم‌زمان، Client درخواست را ارسال می‌کند و تا کامل‌شدن پردازش منتظر می‌ماند:

POST /generate
انتظار برای پردازش
دریافت نتیجه نهایی

مزایای این روش:

  • پیاده‌سازی ساده
  • مناسب برای پاسخ‌های سریع
  • عدم نیاز به صف یا Worker
  • دریافت مستقیم نتیجه در همان اتصال

محدودیت‌های آن:

  • احتمال Timeout
  • اشغال اتصال HTTP
  • افزایش مصرف Connection و Worker وب
  • تجربه نامناسب برای عملیات طولانی
  • دشوارشدن Retry
  • مقیاس‌پذیری محدود
  • ازبین‌رفتن درخواست در صورت Restart سرور

پردازش غیرهم‌زمان یا Asynchronous Job Processing

در معماری غیرهم‌زمان، API فقط درخواست را ثبت می‌کند:

POST /jobs
→ 202 Accepted
→ job_id

سپس Client وضعیت درخواست را جداگانه دریافت می‌کند:

GET /jobs/{job_id}

یا Backend پس از پایان پردازش، نتیجه را به Webhook اعلام می‌کند.

مزایای این روش:

  • پاسخ سریع API
  • کاهش احتمال Timeout
  • امکان اجرای عملیات طولانی
  • مقیاس‌پذیری مستقل Workerها
  • مدیریت Retry
  • کنترل تعداد پردازش‌های هم‌زمان
  • اولویت‌بندی Jobها
  • نگهداری وضعیت هر پردازش
  • امکان توقف یا لغو Job
  • مناسب برای پردازش‌های دسته‌ای

چه زمانی به Job Queue نیاز داریم؟

استفاده از صف برای تمام درخواست‌ها ضروری نیست. اگر پاسخ مدل معمولاً در چند ثانیه تولید می‌شود و کاربر باید خروجی را به‌صورت Streaming ببیند، معماری هم‌زمان ممکن است انتخاب مناسب‌تری باشد.

از Job Queue زمانی استفاده کنید که حداقل یکی از شرایط زیر برقرار باشد:

  • عملیات ممکن است بیشتر از Timeout معمول HTTP طول بکشد.
  • تعداد درخواست‌ها نوسان زیادی دارد.
  • پردازش به چند مرحله وابسته است.
  • باید خطاهای موقت را Retry کنید.
  • باید هم‌زمانی را کنترل کنید.
  • پردازش برای کاربر فوری نیست.
  • نتیجه می‌تواند بعداً تحویل داده شود.
  • هر Job منابع زیادی مصرف می‌کند.
  • نیاز به اولویت‌بندی کاربران یا وظایف دارید.
  • باید گزارش وضعیت یا درصد پیشرفت نمایش دهید.

برای مثال، یک چت‌بات تعاملی معمولاً از Streaming استفاده می‌کند، اما تولید یک گزارش ۳۰ صفحه‌ای، پردازش ۵۰۰ فایل یا ساخت مجموعه‌ای از تصاویر بهتر است در قالب Job اجرا شود.

معماری پیشنهادی

اجزای اصلی سیستم عبارت‌اند از:

  1. Client درخواست را به FastAPI ارسال می‌کند.
  2. FastAPI ورودی را اعتبارسنجی می‌کند.
  3. یک شناسه منحصربه‌فرد برای Job ساخته می‌شود.
  4. پیام پردازش در Redis قرار می‌گیرد.
  5. FastAPI پاسخ 202 Accepted برمی‌گرداند.
  6. یک Celery Worker پیام را از صف دریافت می‌کند.
  7. Worker درخواست را به API درواره ارسال می‌کند.
  8. نتیجه در Result Backend یا پایگاه داده ذخیره می‌شود.
  9. Client با Polling وضعیت را بررسی می‌کند.
  10. در صورت وجود Callback URL، نتیجه از طریق Webhook ارسال می‌شود.

نکته مهم این است که API وب و Worker دو پردازش مستقل هستند. بنابراین می‌توان تعداد Workerها را بدون افزایش تعداد سرورهای API تغییر داد.

وضعیت‌های استاندارد Job

برای یک سیستم واقعی بهتر است وضعیت‌های Job به‌صورت شفاف تعریف شوند.

وضعیتمفهوم
queuedدرخواست ثبت شده و در انتظار Worker است
processingWorker پردازش را آغاز کرده است
retryingپردازش به دلیل خطای موقت دوباره اجرا می‌شود
succeededJob با موفقیت کامل شده است
failedپردازش پس از تلاش‌های مجاز شکست خورده است
cancelledJob پیش از پایان لغو شده است
expiredنتیجه یا Job از زمان نگهداری مجاز عبور کرده است

Celery از نام‌هایی مانند PENDING، STARTED، RETRY، SUCCESS و FAILURE استفاده می‌کند. در API عمومی محصول می‌توانید این وضعیت‌ها را به نام‌های خواناتر تبدیل کنید.

چرا FastAPI BackgroundTasks کافی نیست؟

FastAPI قابلیتی به نام BackgroundTasks دارد که اجرای یک تابع کوچک را پس از ارسال پاسخ ممکن می‌کند. این قابلیت برای کارهایی مانند ثبت Log یا ارسال یک اعلان ساده مناسب است.

اما برای پردازش‌های سنگین و طولانی محدودیت دارد:

  • Task در همان فرایند برنامه اجرا می‌شود.
  • صف پایدار و مستقل ندارد.
  • مقیاس‌پذیری آن محدود است.
  • Restart شدن برنامه می‌تواند Task را متوقف کند.
  • کنترل پیشرفته Retry و اولویت‌بندی ندارد.
  • اجرای Task روی چند سرور نیازمند معماری دیگری است.

مستندات رسمی FastAPI نیز برای محاسبات سنگین پس‌زمینه، استفاده از ابزارهایی مانند Celery و یک Message Queue مانند Redis یا RabbitMQ را پیشنهاد می‌کند. برای جزئیات بیشتر می‌توانید راهنمای Background Tasks در FastAPI را مطالعه کنید.

قاعده عملی:

  • کار کوچک و کم‌اهمیت: BackgroundTasks
  • پردازش طولانی یا حیاتی: Job Queue و Worker

فناوری‌های این پروژه

در این آموزش از پایتون (Python) و ابزارهای زیر استفاده می‌کنیم:

  • Python 3.11 یا جدیدتر
  • FastAPI
  • Uvicorn
  • Celery
  • Redis
  • HTTPX
  • Pydantic Settings
  • Docker و Docker Compose
  • API هوش مصنوعی درواره

ساختار پروژه:

darvareh-async-api/
├── app/
│   ├── __init__.py
│   ├── config.py
│   ├── celery_app.py
│   ├── tasks.py
│   └── main.py
├── requirements.txt
├── Dockerfile
├── compose.yaml
└── .env

مرحله اول: ساخت پروژه

پوشه پروژه را ایجاد کنید:

mkdir darvareh-async-api
cd darvareh-async-api
mkdir app
touch app/__init__.py

یک محیط مجازی بسازید:

python -m venv .venv

فعال‌سازی در Linux و macOS:

source .venv/bin/activate

فعال‌سازی در Windows PowerShell:

.venv\Scripts\Activate.ps1

مرحله دوم: نصب وابستگی‌ها

فایل requirements.txt را ایجاد کنید:

fastapi
uvicorn[standard]
celery[redis]
redis
httpx
pydantic-settings

سپس پکیج‌ها را نصب کنید:

pip install -r requirements.txt

برای یک محیط Production بهتر است نسخه دقیق وابستگی‌ها را پس از آزمایش پروژه قفل کنید.

مرحله سوم: تنظیم متغیرهای محیطی

فایل .env را ایجاد کنید:

DARVAREH_API_KEY=YOUR_DARVAREH_API_KEY
DARVAREH_BASE_URL=https://api.darvareh.ir/v1
DARVAREH_MODEL_ID=YOUR_MODEL_ID

REDIS_URL=redis://redis:6379/0
RESULT_EXPIRES_SECONDS=86400
WEBHOOK_SECRET=CHANGE_THIS_TO_A_LONG_RANDOM_SECRET

مقدار YOUR_MODEL_ID باید با Model ID واقعی درواره جایگزین شود. شناسه مدل را حدس نزنید؛ آن را از صفحه مدل‌های درواره کپی کنید.

برای مشاهده مدل‌های موجود و قیمت به‌روز آن‌ها، صفحه مدل‌های درواره را ببینید.

کلید API را در کد، مخزن Git، اپلیکیشن موبایل یا JavaScript مرورگر قرار ندهید. این مقدار باید فقط در محیط امن سمت سرور نگهداری شود.

مرحله چهارم: ساخت تنظیمات برنامه

فایل app/config.py:

from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    darvareh_api_key: str
    darvareh_base_url: str = "https://api.darvareh.ir/v1"
    darvareh_model_id: str

    redis_url: str = "redis://localhost:6379/0"
    result_expires_seconds: int = 86400
    webhook_secret: str

    model_config = SettingsConfigDict(
        env_file=".env",
        case_sensitive=False,
        extra="ignore",
    )


settings = Settings()

استفاده از متغیر محیطی باعث می‌شود تنظیمات محیط توسعه، آزمایش و Production از کد برنامه جدا بمانند.

مرحله پنجم: پیکربندی Celery

فایل app/celery_app.py:

from celery import Celery

from app.config import settings


celery_app = Celery(
    "darvareh_async_api",
    broker=settings.redis_url,
    backend=settings.redis_url,
    include=["app.tasks"],
)

celery_app.conf.update(
    task_serializer="json",
    result_serializer="json",
    accept_content=["json"],
    result_expires=settings.result_expires_seconds,
    task_track_started=True,
    task_acks_late=True,
    worker_prefetch_multiplier=1,
    broker_connection_retry_on_startup=True,
    timezone="UTC",
    enable_utc=True,
)

گزینه‌های مهم:

task_track_started

با فعال‌کردن این گزینه، وضعیت STARTED پس از شروع پردازش ثبت می‌شود. این وضعیت برای Jobهای طولانی کاربردی است.

task_acks_late

Worker پیام را پس از پایان Task تأیید می‌کند. اگر Worker قبل از پایان متوقف شود، Broker می‌تواند پیام را دوباره تحویل دهد.

این رفتار قابلیت اطمینان را افزایش می‌دهد، اما یک پیام ممکن است بیش از یک بار اجرا شود. به همین دلیل Task باید تا حد ممکن Idempotent باشد.

worker_prefetch_multiplier=1

هر Worker تعداد کمی Job را پیشاپیش رزرو می‌کند. این تنظیم برای وظایف طولانی باعث توزیع متعادل‌تر Jobها میان Workerها می‌شود.

مستندات Celery تأکید می‌کند که Taskها بهتر است Idempotent باشند و عملیات شبکه‌ای نیز Timeout مشخص داشته باشند. جزئیات بیشتر در مستندات رسمی Celery Tasks آمده است.

مرحله ششم: ساخت Worker و اتصال به API درواره

فایل app/tasks.py:

import hashlib
import hmac
import json

import httpx

from app.celery_app import celery_app
from app.config import settings


class RetryableAIError(Exception):
    pass


def create_signature(payload: dict) -> str:
    body = json.dumps(
        payload,
        ensure_ascii=False,
        separators=(",", ":"),
        sort_keys=True,
    ).encode("utf-8")

    return hmac.new(
        settings.webhook_secret.encode("utf-8"),
        body,
        hashlib.sha256,
    ).hexdigest()


@celery_app.task(
    bind=True,
    autoretry_for=(httpx.TransportError, RetryableAIError),
    retry_backoff=True,
    retry_backoff_max=120,
    retry_jitter=True,
    max_retries=4,
    acks_late=True,
)
def generate_ai_response(
    self,
    prompt: str,
    callback_url: str | None = None,
):
    request_payload = {
        "model": settings.darvareh_model_id,
        "messages": [
            {
                "role": "system",
                "content": (
                    "پاسخی دقیق، منظم و کاربردی به زبان فارسی ارائه کن."
                ),
            },
            {
                "role": "user",
                "content": prompt,
            },
        ],
        "temperature": 0.3,
    }

    headers = {
        "Authorization": f"Bearer {settings.darvareh_api_key}",
        "Content-Type": "application/json",
    }

    timeout = httpx.Timeout(
        connect=10.0,
        read=120.0,
        write=30.0,
        pool=10.0,
    )

    with httpx.Client(timeout=timeout) as client:
        response = client.post(
            f"{settings.darvareh_base_url}/chat/completions",
            headers=headers,
            json=request_payload,
        )

    if response.status_code == 429 or response.status_code >= 500:
        raise RetryableAIError(
            f"Temporary AI API error: {response.status_code}"
        )

    response.raise_for_status()
    data = response.json()

    result = {
        "job_id": self.request.id,
        "model": data.get("model", settings.darvareh_model_id),
        "content": data["choices"][0]["message"]["content"],
        "usage": data.get("usage"),
    }

    if callback_url:
        deliver_webhook.delay(
            callback_url=callback_url,
            event="job.succeeded",
            job_id=self.request.id,
            result=result,
        )

    return result


@celery_app.task(
    bind=True,
    autoretry_for=(httpx.TransportError, RetryableAIError),
    retry_backoff=True,
    retry_backoff_max=300,
    retry_jitter=True,
    max_retries=6,
)
def deliver_webhook(
    self,
    callback_url: str,
    event: str,
    job_id: str,
    result: dict,
):
    payload = {
        "event": event,
        "job_id": job_id,
        "result": result,
    }

    signature = create_signature(payload)

    headers = {
        "Content-Type": "application/json",
        "X-Darvareh-Job-Signature": f"sha256={signature}",
    }

    with httpx.Client(timeout=15.0) as client:
        response = client.post(
            callback_url,
            headers=headers,
            json=payload,
        )

    if response.status_code == 429 or response.status_code >= 500:
        raise RetryableAIError(
            f"Temporary webhook error: {response.status_code}"
        )

    response.raise_for_status()

    return {
        "delivered": True,
        "status_code": response.status_code,
    }

چرا ارسال Webhook یک Task جداگانه است؟

اگر فراخوانی مدل و ارسال Webhook را در یک Task قرار دهیم، شکست Webhook ممکن است باعث اجرای دوباره کل Task و ارسال مجدد درخواست به مدل شود. نتیجه این اتفاق می‌تواند مصرف و هزینه تکراری باشد.

در این مثال:

  • generate_ai_response فقط مدل را فراخوانی می‌کند.
  • deliver_webhook نتیجه آماده‌شده را تحویل می‌دهد.
  • شکست Webhook باعث تکرار فراخوانی مدل نمی‌شود.
  • سیاست Retry هر بخش مستقل است.

این جداسازی یکی از مهم‌ترین اصول طراحی گردش‌کارهای غیرهم‌زمان است.

چه خطاهایی را باید Retry کنیم؟

تمام خطاها نباید دوباره اجرا شوند.

خطاهای مناسب برای Retry:

  • خطای شبکه
  • قطع موقت اتصال
  • HTTP 429
  • HTTP 500
  • HTTP 502
  • HTTP 503
  • HTTP 504
  • Timeout موقت

خطاهایی که معمولاً نباید Retry شوند:

  • API Key نامعتبر
  • Model ID اشتباه
  • ورودی نامعتبر
  • Payload ناسازگار
  • درخواست ممنوع
  • خطای اعتبارسنجی
  • فایل با فرمت پشتیبانی‌نشده

Retry کردن خطای دائمی فقط صف را شلوغ و منابع را مصرف می‌کند.

در Celery می‌توان از Exponential Backoff و Jitter استفاده کرد. طبق مستندات رسمی Retry در Celery، Backoff فاصله میان تلاش‌ها را تدریجی افزایش می‌دهد و Jitter از اجرای هم‌زمان تعداد زیادی Retry جلوگیری می‌کند.

مرحله هفتم: ساخت REST API با FastAPI

فایل app/main.py:

from typing import Any
from urllib.parse import urlparse
from uuid import uuid4

import redis
from celery.result import AsyncResult
from fastapi import FastAPI, Header, HTTPException, Response, status
from pydantic import AnyHttpUrl, BaseModel, Field

from app.celery_app import celery_app
from app.config import settings
from app.tasks import generate_ai_response


app = FastAPI(
    title="Darvareh Async AI API",
    version="1.0.0",
)

redis_client = redis.Redis.from_url(
    settings.redis_url,
    decode_responses=True,
)


class CreateJobRequest(BaseModel):
    prompt: str = Field(min_length=1, max_length=20_000)
    callback_url: AnyHttpUrl | None = None


class CreateJobResponse(BaseModel):
    job_id: str
    status: str
    status_url: str
    duplicate: bool = False


class JobStatusResponse(BaseModel):
    job_id: str
    status: str
    result: Any | None = None
    error: str | None = None


def validate_callback_url(callback_url: str | None) -> None:
    if not callback_url:
        return

    parsed = urlparse(callback_url)

    if parsed.scheme != "https":
        raise HTTPException(
            status_code=422,
            detail="callback_url must use HTTPS",
        )

    if parsed.hostname in {"localhost", "127.0.0.1", "::1"}:
        raise HTTPException(
            status_code=422,
            detail="Local callback URLs are not allowed",
        )


def map_celery_state(state: str) -> str:
    states = {
        "PENDING": "queued",
        "STARTED": "processing",
        "RETRY": "retrying",
        "SUCCESS": "succeeded",
        "FAILURE": "failed",
        "REVOKED": "cancelled",
    }

    return states.get(state, state.lower())


@app.get("/health")
def health():
    try:
        redis_client.ping()
    except redis.RedisError as exc:
        raise HTTPException(
            status_code=503,
            detail="Redis is unavailable",
        ) from exc

    return {"status": "ok"}


@app.post(
    "/jobs",
    response_model=CreateJobResponse,
    status_code=status.HTTP_202_ACCEPTED,
)
def create_job(
    body: CreateJobRequest,
    response: Response,
    idempotency_key: str | None = Header(
        default=None,
        alias="Idempotency-Key",
    ),
):
    callback_url = (
        str(body.callback_url) if body.callback_url else None
    )

    validate_callback_url(callback_url)

    if idempotency_key:
        if len(idempotency_key) > 200:
            raise HTTPException(
                status_code=422,
                detail="Idempotency-Key is too long",
            )

        redis_key = f"idempotency:{idempotency_key}"
        existing_job_id = redis_client.get(redis_key)

        if existing_job_id:
            return CreateJobResponse(
                job_id=existing_job_id,
                status="queued",
                status_url=f"/jobs/{existing_job_id}",
                duplicate=True,
            )
    else:
        redis_key = None

    job_id = str(uuid4())

    if redis_key:
        created = redis_client.set(
            redis_key,
            job_id,
            nx=True,
            ex=settings.result_expires_seconds,
        )

        if not created:
            existing_job_id = redis_client.get(redis_key)

            return CreateJobResponse(
                job_id=existing_job_id,
                status="queued",
                status_url=f"/jobs/{existing_job_id}",
                duplicate=True,
            )

    try:
        generate_ai_response.apply_async(
            task_id=job_id,
            kwargs={
                "prompt": body.prompt,
                "callback_url": callback_url,
            },
        )
    except Exception:
        if redis_key:
            redis_client.delete(redis_key)
        raise

    response.headers["Location"] = f"/jobs/{job_id}"

    return CreateJobResponse(
        job_id=job_id,
        status="queued",
        status_url=f"/jobs/{job_id}",
    )


@app.get(
    "/jobs/{job_id}",
    response_model=JobStatusResponse,
)
def get_job(job_id: str):
    task_result = AsyncResult(
        job_id,
        app=celery_app,
    )

    public_status = map_celery_state(task_result.state)

    if task_result.state == "SUCCESS":
        return JobStatusResponse(
            job_id=job_id,
            status=public_status,
            result=task_result.result,
        )

    if task_result.state == "FAILURE":
        return JobStatusResponse(
            job_id=job_id,
            status=public_status,
            error=str(task_result.result),
        )

    if task_result.state == "RETRY":
        return JobStatusResponse(
            job_id=job_id,
            status=public_status,
            error="Temporary error; the job will be retried",
        )

    return JobStatusResponse(
        job_id=job_id,
        status=public_status,
    )


@app.delete(
    "/jobs/{job_id}",
    status_code=status.HTTP_202_ACCEPTED,
)
def cancel_job(job_id: str):
    celery_app.control.revoke(
        job_id,
        terminate=False,
    )

    return {
        "job_id": job_id,
        "status": "cancellation_requested",
    }

Idempotency Key چیست؟

فرض کنید Client درخواست ساخت یک گزارش را ارسال می‌کند، اما پیش از دریافت پاسخ اتصال قطع می‌شود. Client نمی‌داند درخواست ثبت شده یا نه؛ بنابراین آن را دوباره ارسال می‌کند.

بدون Idempotency، دو Job ساخته می‌شود و مدل دو بار فراخوانی خواهد شد.

Client می‌تواند یک کلید ثابت همراه درخواست بفرستد:

Idempotency-Key: report-order-825-user-42

اگر همان درخواست دوباره ارسال شود، API شناسه Job قبلی را برمی‌گرداند.

نمونه درخواست:

curl -X POST http://localhost:8000/jobs \
  -H "Content-Type: application/json" \
  -H "Idempotency-Key: customer-report-1001" \
  -d '{
    "prompt": "یک گزارش مدیریتی درباره بازخورد مشتریان بنویس."
  }'

نکته مهم این است که Idempotency Key باید یک عملیات منطقی را نمایش دهد. استفاده از یک کلید ثابت برای تمام درخواست‌های کاربر اشتباه است.

در محیط Production بهتر است علاوه بر کلید، Hash ورودی نیز ذخیره شود. اگر Client همان کلید را با Payload متفاوت ارسال کرد، API باید خطای Conflict برگرداند.

مرحله هشتم: ساخت Dockerfile

فایل Dockerfile:

FROM python:3.12-slim

ENV PYTHONDONTWRITEBYTECODE=1
ENV PYTHONUNBUFFERED=1

WORKDIR /app

COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY app ./app

CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]

مرحله نهم: اجرای سرویس‌ها با Docker Compose

فایل compose.yaml:

services:
  redis:
    image: redis:7-alpine
    restart: unless-stopped
    command:
      - redis-server
      - --appendonly
      - "yes"
    volumes:
      - redis-data:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 5s
      timeout: 3s
      retries: 10

  api:
    build: .
    restart: unless-stopped
    env_file:
      - .env
    ports:
      - "8000:8000"
    depends_on:
      redis:
        condition: service_healthy

  worker:
    build: .
    restart: unless-stopped
    command:
      - celery
      - -A
      - app.celery_app.celery_app
      - worker
      - --loglevel=INFO
      - --concurrency=2
    env_file:
      - .env
    depends_on:
      redis:
        condition: service_healthy

volumes:
  redis-data:

اجرای پروژه:

docker compose up --build

پس از اجرا، مستندات Swagger در آدرس زیر در دسترس است:

http://localhost:8000/docs

بررسی سلامت سرویس:

curl http://localhost:8000/health

پاسخ مورد انتظار:

{
  "status": "ok"
}

مرحله دهم: ثبت اولین Job

درخواست زیر را ارسال کنید:

curl -X POST http://localhost:8000/jobs \
  -H "Content-Type: application/json" \
  -H "Idempotency-Key: article-summary-001" \
  -d '{
    "prompt": "پنج کاربرد عملی هوش مصنوعی در فروشگاه اینترنتی را توضیح بده."
  }'

پاسخ:

{
  "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
  "status": "queued",
  "status_url": "/jobs/0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
  "duplicate": false
}

کد وضعیت HTTP برابر 202 Accepted است. این کد به معنی پذیرفته‌شدن درخواست برای پردازش است، نه کامل‌شدن آن.

مرحله یازدهم: دریافت نتیجه با Polling

وضعیت Job را بررسی کنید:

curl \
  http://localhost:8000/jobs/0a6c2b26-84e4-4c59-8a72-7fd1166f7591

هنگام انتظار:

{
  "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
  "status": "queued",
  "result": null,
  "error": null
}

هنگام اجرا:

{
  "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
  "status": "processing",
  "result": null,
  "error": null
}

پس از موفقیت:

{
  "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
  "status": "succeeded",
  "result": {
    "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
    "model": "YOUR_MODEL_ID",
    "content": "پاسخ تولیدشده توسط مدل...",
    "usage": {
      "prompt_tokens": 42,
      "completion_tokens": 380,
      "total_tokens": 422
    }
  },
  "error": null
}

مقادیر دقیق بخش usage به مدل انتخاب‌شده و پاسخ API بستگی دارد.

پیاده‌سازی Polling در JavaScript

نمونه ساده برای Frontend:

async function waitForJob(jobId) {
  let delay = 1000;

  while (true) {
    const response = await fetch(`/jobs/${jobId}`);

    if (!response.ok) {
      throw new Error("دریافت وضعیت Job ناموفق بود");
    }

    const job = await response.json();

    if (job.status === "succeeded") {
      return job.result;
    }

    if (
      job.status === "failed" ||
      job.status === "cancelled" ||
      job.status === "expired"
    ) {
      throw new Error(job.error || "پردازش ناموفق بود");
    }

    await new Promise((resolve) => {
      setTimeout(resolve, delay);
    });

    delay = Math.min(delay * 1.5, 10_000);
  }
}

در این مثال فاصله Polling تدریجی افزایش پیدا می‌کند:

1s → 1.5s → 2.25s → 3.37s → ... → حداکثر 10s

ارسال درخواست وضعیت در هر ۱۰۰ میلی‌ثانیه باعث ایجاد بار غیرضروری روی Backend و Redis می‌شود. برای بیشتر پروژه‌ها، فاصله اولیه یک تا دو ثانیه مناسب است.

استفاده از Webhook

Polling برای رابط کاربری ساده است، اما اگر Backend دیگری منتظر نتیجه باشد، Webhook کارآمدتر خواهد بود.

ثبت Job همراه Callback URL:

curl -X POST http://localhost:8000/jobs \
  -H "Content-Type: application/json" \
  -H "Idempotency-Key: report-with-webhook-001" \
  -d '{
    "prompt": "یک گزارش کوتاه از روندهای کاربردی هوش مصنوعی تهیه کن.",
    "callback_url": "https://example.com/webhooks/ai-jobs"
  }'

پس از پایان Job، سرویس چنین Payloadی ارسال می‌کند:

{
  "event": "job.succeeded",
  "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
  "result": {
    "job_id": "0a6c2b26-84e4-4c59-8a72-7fd1166f7591",
    "model": "YOUR_MODEL_ID",
    "content": "نتیجه پردازش...",
    "usage": {
      "prompt_tokens": 35,
      "completion_tokens": 290,
      "total_tokens": 325
    }
  }
}

امضای پیام نیز در Header قرار می‌گیرد:

X-Darvareh-Job-Signature: sha256=...

نام این Header در نمونه، مربوط به سرویس Job خودمان است و یک Header رسمی API درواره محسوب نمی‌شود.

اعتبارسنجی امضای Webhook

نمونه دریافت Webhook با FastAPI:

import hashlib
import hmac

from fastapi import FastAPI, Header, HTTPException, Request


app = FastAPI()

WEBHOOK_SECRET = "CHANGE_THIS_TO_THE_SAME_SECRET"


@app.post("/webhooks/ai-jobs")
async def receive_ai_job(
    request: Request,
    x_darvareh_job_signature: str = Header(),
):
    raw_body = await request.body()

    expected = hmac.new(
        WEBHOOK_SECRET.encode("utf-8"),
        raw_body,
        hashlib.sha256,
    ).hexdigest()

    received = x_darvareh_job_signature.removeprefix(
        "sha256="
    )

    if not hmac.compare_digest(expected, received):
        raise HTTPException(
            status_code=401,
            detail="Invalid webhook signature",
        )

    payload = await request.json()

    return {
        "received": True,
        "job_id": payload["job_id"],
    }

دریافت‌کننده Webhook باید سریع پاسخ دهد. اگر پردازش دیگری لازم است، خود دریافت‌کننده نیز بهتر است پیام را در صف داخلی قرار دهد و پاسخ 2xx برگرداند.

Polling یا Webhook؛ کدام بهتر است؟

معیارPollingWebhook
پیاده‌سازی Frontendسادهمعمولاً به Backend نیاز دارد
مصرف درخواستبیشترکمتر
دریافت نزدیک به لحظه نتیجهوابسته به فاصله Pollingبله
مناسب Browserبلهمستقیم خیر
مناسب ارتباط سرور با سرورقابل‌استفادهمناسب‌تر
نیاز به Endpoint عمومیخیربله
مدیریت Retry تحویلسمت Clientسمت فرستنده Webhook
پیچیدگی اعتبارسنجیکمتربیشتر

معماری پیشنهادی برای بسیاری از محصولات، پشتیبانی هم‌زمان از هر دو روش است:

  • Frontend از Polling استفاده کند.
  • سرویس‌های Backend از Webhook استفاده کنند.
  • وضعیت نهایی همیشه از GET /jobs/{id} قابل‌بازیابی باشد.

Webhook باید نقش اعلان را داشته باشد، نه تنها محل نگهداری نتیجه. ممکن است دریافت‌کننده هنگام ارسال Webhook موقتاً در دسترس نباشد.

ذخیره Jobها در PostgreSQL

استفاده از Redis Result Backend برای نمونه آموزشی مناسب است، اما در یک محصول واقعی معمولاً وضعیت پایدار Jobها در PostgreSQL ذخیره می‌شود.

ساختار پیشنهادی جدول:

CREATE TABLE ai_jobs (
    id UUID PRIMARY KEY,
    user_id UUID,
    idempotency_key VARCHAR(200),
    job_type VARCHAR(100) NOT NULL,
    status VARCHAR(30) NOT NULL,
    model_id VARCHAR(200),
    input_payload JSONB NOT NULL,
    result_payload JSONB,
    error_code VARCHAR(100),
    error_message TEXT,
    progress SMALLINT DEFAULT 0,
    attempt_count INTEGER DEFAULT 0,
    callback_url TEXT,
    created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    started_at TIMESTAMPTZ,
    completed_at TIMESTAMPTZ,
    expires_at TIMESTAMPTZ
);

قید یکتایی Idempotency باید در محدوده کاربر تعریف شود:

CREATE UNIQUE INDEX ai_jobs_user_idempotency_unique
ON ai_jobs (user_id, idempotency_key)
WHERE idempotency_key IS NOT NULL;

ایندکس وضعیت و زمان ایجاد:

CREATE INDEX ai_jobs_status_created_at_idx
ON ai_jobs (status, created_at);

مزایای PostgreSQL:

  • نگهداری تاریخچه Job
  • گزارش‌گیری
  • محاسبه هزینه
  • مشاهده وضعیت پس از انقضای Redis
  • جست‌وجوی Jobهای ناموفق
  • پیاده‌سازی کنترل دسترسی
  • ذخیره زمان شروع و پایان
  • نگهداری خطاها و تعداد تلاش‌ها

Redis را برای صف، Lock و Cache نگه دارید و PostgreSQL را منبع اصلی وضعیت کسب‌وکار در نظر بگیرید.

نمایش درصد پیشرفت

مدل‌های متنی معمولاً درصد پیشرفت واقعی ارائه نمی‌کنند. نباید درصدی غیرواقعی مانند ۷۳٪ را صرفاً بر اساس زمان نمایش دهید.

اما برای گردش‌کارهای چندمرحله‌ای می‌توان پیشرفت مرحله‌ای تعریف کرد:

10%  اعتبارسنجی ورودی
25%  استخراج متن فایل
45%  بخش‌بندی محتوا
70%  پردازش با مدل
90%  ساخت خروجی
100% تکمیل

در Celery می‌توان وضعیت سفارشی ثبت کرد:

self.update_state(
    state="PROGRESS",
    meta={
        "progress": 45,
        "stage": "processing_chunks",
    },
)

سپس Endpoint وضعیت می‌تواند مقدار task_result.info را برای وضعیت PROGRESS بخواند.

درصد پیشرفت باید بر اساس مراحل واقعی یا تعداد آیتم‌های کامل‌شده محاسبه شود:

progress = completed_items / total_items × 100

این فرمول را می‌توان به شکل زیر در کد نوشت:

progress = int(
    completed_items / total_items * 100
)

مدیریت Visibility Timeout در Redis

وقتی Redis به‌عنوان Broker استفاده می‌شود، Visibility Timeout تعیین می‌کند Worker چه مدت برای تأیید پیام فرصت دارد. اگر پیام در این زمان تأیید نشود، ممکن است دوباره به Worker دیگری تحویل داده شود.

بر اساس مستندات رسمی Celery برای Redis، مقدار پیش‌فرض Visibility Timeout در Redis یک ساعت است.

اگر اجرای یک Job بیشتر از این مقدار طول بکشد، احتمال اجرای مجدد آن وجود دارد.

نمونه تنظیم:

celery_app.conf.broker_transport_options = {
    "visibility_timeout": 3600,
}

افزایش بدون بررسی این مقدار نیز راه‌حل کاملی نیست؛ زیرا در صورت ازبین‌رفتن Worker، تحویل مجدد Job به تأخیر می‌افتد.

راهکارهای بهتر:

  • Jobهای بسیار طولانی را به مراحل کوچک‌تر تقسیم کنید.
  • هر مرحله را Idempotent طراحی کنید.
  • زمان اجرای واقعی را اندازه‌گیری کنید.
  • برای Taskهای کوتاه و بلند صف جدا بسازید.
  • از Timeout شبکه استفاده کنید.
  • نتیجه هر مرحله را ذخیره کنید.
  • Workerهای مخصوص پردازش طولانی داشته باشید.

مفهوم At-Least-Once Delivery

بسیاری از سیستم‌های صف تلاش می‌کنند هر پیام «حداقل یک بار» اجرا شود. این سیاست بهتر از گم‌شدن Job است، اما اجرای تکراری را کاملاً حذف نمی‌کند.

در نتیجه باید فرض کنید یک Task ممکن است دو بار اجرا شود.

عملیات زیر در صورت تکرار می‌توانند مشکل ایجاد کنند:

  • کسر دوباره اعتبار
  • ارسال چندباره پیام
  • ثبت سفارش تکراری
  • ایجاد چند فایل یکسان
  • فراخوانی تکراری مدل و افزایش هزینه
  • ارسال چندباره Webhook

راهکارها:

  • استفاده از Idempotency Key
  • ذخیره شناسه خارجی عملیات
  • قید یکتا در پایگاه داده
  • بررسی وضعیت پیش از اجرای اثر جانبی
  • ثبت Ledger برای هزینه
  • تفکیک تولید نتیجه از تحویل نتیجه
  • استفاده از Transaction
  • ذخیره Hash ورودی و خروجی

صف‌های جداگانه برای وظایف مختلف

قرار دادن تمام Jobها در یک صف می‌تواند باعث شود یک ویدئوی طولانی، صدها درخواست متنی کوتاه را پشت صف نگه دارد.

پیشنهاد:

ai_text
ai_document
ai_image
ai_audio
ai_video
webhooks

مسیریابی در Celery:

celery_app.conf.task_routes = {
    "app.tasks.generate_ai_response": {
        "queue": "ai_text",
    },
    "app.tasks.deliver_webhook": {
        "queue": "webhooks",
    },
}

اجرای Worker متنی:

celery -A app.celery_app.celery_app worker \
  --loglevel=INFO \
  --queues=ai_text \
  --concurrency=4

اجرای Worker مخصوص Webhook:

celery -A app.celery_app.celery_app worker \
  --loglevel=INFO \
  --queues=webhooks \
  --concurrency=8

به این ترتیب می‌توانید هر گروه را مستقل مقیاس دهید.

اولویت‌بندی Jobها

همه درخواست‌ها ارزش و فوریت یکسانی ندارند. برای مثال:

  • درخواست تعاملی کاربر باید سریع اجرا شود.
  • پردازش Batch شبانه می‌تواند صبر کند.
  • کاربران سازمانی ممکن است صف اختصاصی داشته باشند.
  • Webhookهای نتیجه نباید پشت تولیدهای سنگین بمانند.

رویکرد قابل‌فهم‌تر از یک صف بسیار پیچیده، ساخت صف‌های مجزا است:

ai_interactive
ai_default
ai_batch

سپس بر اساس نوع درخواست، Job را به صف مناسب بفرستید:

generate_ai_response.apply_async(
    kwargs={
        "prompt": body.prompt,
        "callback_url": callback_url,
    },
    queue="ai_interactive",
)

کنترل هم‌زمانی و هزینه

افزایش تعداد Workerها همیشه باعث عملکرد بهتر نمی‌شود. هم‌زمانی زیاد می‌تواند پیامدهای زیر را داشته باشد:

  • رسیدن سریع‌تر به Rate Limit
  • افزایش Retry
  • افزایش مصرف API
  • اشباع اتصال‌های شبکه
  • فشار روی Redis
  • افزایش هزینه ناگهانی
  • کاهش پایداری سیستم

یک تخمین ساده برای ظرفیت:

throughput ≈ concurrency / average_job_duration

اگر میانگین زمان هر Job برابر ۲۰ ثانیه و Concurrency برابر ۱۰ باشد:

throughput ≈ 10 / 20 = 0.5 job per second

یعنی حدود ۳۰ Job در دقیقه، البته در شرایط ایدئال.

برای کنترل بهتر:

  • Concurrency را به‌تدریج افزایش دهید.
  • Queue Length را مانیتور کنید.
  • زمان انتظار در صف را اندازه بگیرید.
  • خطاهای 429 را ثبت کنید.
  • مصرف توکن را برای هر Job نگه دارید.
  • برای هر کاربر سقف مصرف تعریف کنید.
  • درخواست‌های مشابه را Cache کنید.
  • Max Tokens مناسب تعیین کنید.
  • از مدل متناسب با پیچیدگی وظیفه استفاده کنید.

برای مشاهده قیمت و انتخاب مدل مناسب، به فهرست مدل‌ها و قیمت‌های درواره مراجعه کنید.

مانیتورینگ چه شاخص‌هایی ضروری است؟

حداقل این معیارها را ثبت کنید:

معیارهای صف

  • تعداد Jobهای منتظر
  • قدیمی‌ترین Job صف
  • نرخ ورود Job
  • نرخ تکمیل Job
  • نرخ شکست
  • تعداد Retry
  • تعداد Worker فعال

معیارهای زمانی

  • Queue Wait Time
  • Processing Time
  • End-to-End Latency
  • Webhook Delivery Time

تعریف زمان کل:

End-to-End Latency =
Queue Wait Time + Processing Time + Delivery Time

معیارهای API هوش مصنوعی

  • مدل استفاده‌شده
  • تعداد توکن ورودی
  • تعداد توکن خروجی
  • وضعیت HTTP
  • تعداد خطاهای 429
  • تعداد Timeout
  • هزینه تخمینی هر Job
  • نرخ موفقیت هر مدل

معیارهای کسب‌وکار

  • تعداد Job هر کاربر
  • تعداد Job هر قابلیت
  • Jobهای لغوشده
  • Jobهای تکراری
  • هزینه هر مشتری
  • نرخ استفاده از نتیجه

در Log هر رویداد از شناسه‌های زیر استفاده کنید:

request_id
job_id
user_id
task_id
model_id
attempt_number

این شناسه‌ها پیدا کردن مسیر کامل یک درخواست را آسان می‌کنند.

لغو Job چگونه کار می‌کند؟

در مثال از دستور زیر استفاده کردیم:

celery_app.control.revoke(
    job_id,
    terminate=False,
)

اگر Job هنوز اجرا نشده باشد، Worker می‌تواند از اجرای آن صرف‌نظر کند. اما اگر فراخوانی API آغاز شده باشد، لغو Celery لزوماً درخواست شبکه‌ای در حال اجرا را متوقف نمی‌کند.

بنابراین مفهوم لغو باید دقیق تعریف شود:

  • cancellation_requested: درخواست لغو ثبت شده است.
  • cancelled: Job پیش از انجام اثر اصلی متوقف شده است.
  • succeeded: Job پیش از اعمال لغو کامل شده است.
  • cancellation_failed: توقف عملیات ممکن نبوده است.

استفاده از terminate=True می‌تواند فرایند Worker را به‌زور متوقف کند و برای بسیاری از سناریوهای Production انتخاب مناسبی نیست.

برای گردش‌کارهای چندمرحله‌ای، Worker باید بین مراحل وضعیت لغو را از پایگاه داده بررسی کند.

تفاوت Queue، Broker و Result Backend

این سه مفهوم را نباید یکی دانست.

Message Broker

پیام Task را از API به Worker منتقل می‌کند.

نمونه‌ها:

  • Redis
  • RabbitMQ

Worker

پیام را دریافت و کد Task را اجرا می‌کند.

Result Backend

وضعیت و خروجی Task را نگه می‌دارد.

در این آموزش Redis هم Broker و هم Result Backend است. در محصول بزرگ‌تر می‌توانید از Redis برای Broker و PostgreSQL برای وضعیت پایدار Job استفاده کنید.

Broker پایگاه داده کسب‌وکار نیست. اطلاعات مهم مانند مالک Job، هزینه، خروجی نهایی و تاریخچه نباید فقط در صف نگهداری شود.

اتصال همین معماری به تولید تصویر، صوت و ویدئو

الگوی صف به نوع مدل وابسته نیست:

ثبت Job
دریافت توسط Worker
فراخوانی مدل
ذخیره نتیجه
اعلام وضعیت
تحویل خروجی

اما Endpoint و پارامترهای تولید تصویر، صدا یا ویدئو ممکن است بر اساس مدل انتخابی متفاوت باشند. بنابراین Endpoint یا Payload را حدس نزنید و مشخصات مدل فعال را در مستندات درواره بررسی کنید.

برای فایل‌های خروجی بزرگ بهتر است:

  • فایل را در Object Storage ذخیره کنید.
  • در Result فقط URL محدود و زمان‌دار برگردانید.
  • فایل بزرگ را داخل Redis یا جدول Job ذخیره نکنید.
  • زمان انقضای لینک را مشخص کنید.
  • Metadata فایل را در PostgreSQL نگه دارید.
  • پاک‌سازی فایل‌های منقضی را زمان‌بندی کنید.

خطاهای رایج در معماری Async

نگه‌داشتن نتیجه فقط در حافظه API

اگر سرور Restart شود یا درخواست بعدی به Instance دیگری برسد، وضعیت از بین می‌رود.

راه‌حل: استفاده از Redis یا پایگاه داده مشترک.

Retry کردن تمام خطاها

خطاهای دائمی بارها تکرار می‌شوند و صف را اشغال می‌کنند.

راه‌حل: تفکیک خطاهای موقت و دائمی.

اجرای دوباره مدل هنگام شکست Webhook

باعث افزایش هزینه و ایجاد خروجی تکراری می‌شود.

راه‌حل: Task مستقل برای تحویل Webhook.

ذخیره فایل بزرگ در Redis

حافظه Redis به‌سرعت مصرف می‌شود.

راه‌حل: Object Storage و ذخیره URL در نتیجه.

Polling بسیار سریع

Backend و Redis را بی‌دلیل تحت فشار قرار می‌دهد.

راه‌حل: Polling با فاصله افزایشی.

نبود Idempotency

درخواست تکراری باعث اجرای چندباره و هزینه اضافه می‌شود.

راه‌حل: Idempotency Key و قید یکتا.

قرار دادن همه Taskها در یک صف

Jobهای طولانی درخواست‌های سریع را متوقف می‌کنند.

راه‌حل: صف و Worker جدا بر اساس نوع پردازش.

نداشتن Timeout

یک اتصال شبکه‌ای معیوب می‌تواند Worker را مدت زیادی اشغال کند.

راه‌حل: Timeout مستقل برای Connect، Read و Write.

نمایش درصد پیشرفت ساختگی

اعتماد کاربر را کاهش می‌دهد.

راه‌حل: نمایش مرحله واقعی یا وضعیت کلی.

افشای خطای داخلی

نمایش Traceback کامل در API می‌تواند اطلاعات غیرضروری زیرساخت را آشکار کند.

راه‌حل: ذخیره جزئیات در Log و ارائه پیام کنترل‌شده به Client.

مسیر ارتقا از نمونه آموزشی به Production

نمونه این مقاله برای یادگیری و شروع پروژه مناسب است. برای محیط عملیاتی این مراحل را اضافه کنید:

  1. احراز هویت کاربران
  2. ذخیره Job در PostgreSQL
  3. محدودیت تعداد Job برای هر کاربر
  4. محاسبه و ثبت مصرف هر Job
  5. کنترل Callback URL با Allowlist
  6. ذخیره Secretها در Secret Manager
  7. قید یکتای Idempotency در پایگاه داده
  8. صف جدا برای هر نوع Task
  9. Worker جدا برای وظایف کوتاه و بلند
  10. مانیتورینگ Queue Length و Latency
  11. Dashboard مدیریتی Jobها
  12. پاک‌سازی Jobهای منقضی
  13. Object Storage برای فایل‌ها
  14. تست قطع Redis و Worker
  15. تست اجرای تکراری Task
  16. Graceful Shutdown
  17. Health Check جدا برای API و Worker
  18. ثبت Request ID و Job ID
  19. سیاست مشخص Retry و Dead Letter
  20. هشدار برای افزایش خطا یا طول صف

چک‌لیست نهایی

پیش از استقرار سیستم بررسی کنید:

  • درخواست ثبت Job با 202 Accepted پاسخ داده می‌شود.
  • Job ID برای هر درخواست ساخته می‌شود.
  • وضعیت Job قابل‌بازیابی است.
  • کلید API فقط در Worker یا Backend قرار دارد.
  • Timeout شبکه تنظیم شده است.
  • خطاهای 429 و 5xx Retry می‌شوند.
  • خطاهای 4xx دائمی Retry نمی‌شوند.
  • Retry دارای Backoff و Jitter است.
  • Taskهای اثرگذار Idempotent هستند.
  • Webhook Task مستقل دارد.
  • امضای Webhook اعتبارسنجی می‌شود.
  • Polling فاصله منطقی دارد.
  • نتیجه دائمی در پایگاه داده ذخیره می‌شود.
  • فایل‌های بزرگ وارد Redis نمی‌شوند.
  • صف‌های کوتاه و بلند از هم جدا هستند.
  • تعداد Workerها با Rate Limit هماهنگ است.
  • مصرف و هزینه هر Job ثبت می‌شود.
  • Logها شامل Job ID هستند.
  • Jobهای منقضی پاک‌سازی می‌شوند.
  • سناریوی Crash شدن Worker آزمایش شده است.

پرسش‌های متداول

آیا برای هر API هوش مصنوعی به Celery نیاز داریم؟

خیر. برای درخواست‌های سریع و تعاملی، فراخوانی مستقیم یا Streaming ساده‌تر است. Celery برای عملیات طولانی، دسته‌ای، قابل‌تکرار و نیازمند صف مناسب‌تر است.

آیا می‌توان به‌جای Redis از RabbitMQ استفاده کرد؟

بله. RabbitMQ یک Message Broker تخصصی است و برای معماری‌های صف پیچیده انتخاب مناسبی محسوب می‌شود. Redis راه‌اندازی ساده‌تری دارد و برای بسیاری از پروژه‌ها نقطه شروع خوبی است.

آیا Redis برای نگهداری دائمی نتیجه کافی است؟

برای نمونه و Cache بله، اما برای تاریخچه کسب‌وکار بهتر است PostgreSQL یا پایگاه داده پایدار دیگری داشته باشید.

آیا Polling روش بدی است؟

خیر. Polling با فاصله منطقی برای رابط کاربری بسیار قابل‌اعتماد و ساده است. مشکل زمانی ایجاد می‌شود که فاصله درخواست‌ها بیش از حد کوتاه باشد.

آیا Webhook جای Endpoint وضعیت را می‌گیرد؟

خیر. Webhook یک اعلان است و ممکن است تحویل آن با تأخیر یا خطا مواجه شود. نتیجه باید همچنان از Endpoint وضعیت قابل‌بازیابی باشد.

چگونه از پردازش دوباره یک Job جلوگیری کنیم؟

از Idempotency Key، قید یکتا، ثبت وضعیت در پایگاه داده و طراحی Idempotent Worker استفاده کنید. در سیستم صف باید احتمال تحویل تکراری پیام را در نظر بگیرید.

آیا می‌توان تعداد Workerها را خودکار افزایش داد؟

بله. می‌توان بر اساس طول صف، زمان انتظار یا مصرف منابع، Workerها را مقیاس داد. بااین‌حال افزایش Worker باید با Rate Limit و بودجه API هماهنگ باشد.

آیا این معماری برای تولید تصویر و ویدئو مناسب است؟

بله. Job Queue برای عملیات طولانی مانند تولید تصویر، صوت و ویدئو بسیار مناسب است. فقط Endpoint و Payload هر مدل را مطابق مستندات آن مدل تنظیم کنید.

Model ID در کد را از کجا بگیریم؟

مقدار YOUR_MODEL_ID را با Model ID درواره جایگزین کنید. فهرست مدل‌ها و قیمت‌های به‌روز در صفحه مدل‌های درواره قرار دارد.

جمع‌بندی

اتصال مستقیم Backend به مدل هوش مصنوعی برای نمونه اولیه ساده است، اما برای پردازش‌های طولانی، دسته‌ای و پرتعداد به‌تنهایی کافی نیست.

در معماری غیرهم‌زمان:

  • API درخواست را ثبت می‌کند.
  • Job وارد صف می‌شود.
  • Worker پردازش را انجام می‌دهد.
  • وضعیت از طریق Polling قابل‌دریافت است.
  • نتیجه از طریق Webhook نیز اعلام می‌شود.
  • Retry خطاهای موقت را مدیریت می‌کند.
  • Idempotency از پردازش ناخواسته تکراری جلوگیری می‌کند.
  • Workerها مستقل از API مقیاس پیدا می‌کنند.

در این مقاله یک نمونه عملی با FastAPI، Celery و Redis ساختیم و Worker را به API سازگار درواره متصل کردیم. این معماری را می‌توان برای تولید محتوا، تحلیل اسناد، پردازش دسته‌ای، ساخت تصویر، تولید صوت، ویدئو و گردش‌کارهای چندمرحله‌ای توسعه داد.

برای شروع، در درواره ثبت‌نام کنید، کلید API خود را بسازید و Model ID متناسب با کاربرد پروژه را از صفحه مدل‌ها انتخاب کنید.

مقالات مرتبط

برای مطالعه شرایط استفاده و محدودیت‌های مسئولیت، صفحه «سلب مسئولیت» را مشاهده کنید.

Read more