RabbitMQ چیست؟ آموزش کامل Message Queue با Python، FastAPI و API درواره

RabbitMQ یک Message Broker برای اجرای Jobهای غیرهم‌زمان و ارتباط میان سرویس‌هاست. در این آموزش، RabbitMQ را با Docker اجرا می‌کنید و یک سیستم عملی با Python، FastAPI، Retry، DLQ و API درواره می‌سازید.

Share
RabbitMQ چیست؟ آموزش کامل Message Queue با Python، FastAPI و API درواره

در یک برنامه ساده، معمولاً همه عملیات داخل همان HTTP Request اجرا می‌شوند. کاربر درخواست را ارسال می‌کند، Backend عملیات لازم را انجام می‌دهد و پس از پایان پردازش پاسخ را برمی‌گرداند.

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

اجرای عملیات سنگین داخل Request اصلی مشکلاتی ایجاد می‌کند:

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

RabbitMQ یا ربیت ام کیو یک Message Broker است که پیام‌ها را از Producer دریافت و از طریق Queue به Consumerها تحویل می‌دهد. این ابزار به شما اجازه می‌دهد عملیات سنگین را از مسیر اصلی API جدا کنید و Workerهای مستقل برای پردازش آن‌ها بسازید.

در این آموزش ابتدا با مفاهیم RabbitMQ مانند Queue، Exchange، Binding، Routing Key، Acknowledgement و Dead Letter Exchange آشنا می‌شویم. سپس RabbitMQ را با Docker اجرا می‌کنیم و یک سیستم عملی با Python، FastAPI، aio-pika و API هوش مصنوعی درواره می‌سازیم.

RabbitMQ چیست؟

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

در ساده‌ترین حالت، معماری RabbitMQ چنین است:

Producer → RabbitMQ → Queue → Consumer

در پروژه این مقاله:

  • FastAPI نقش Producer را دارد.
  • RabbitMQ پیام‌ها را Route و در Queue نگهداری می‌کند.
  • Worker نقش Consumer را دارد.
  • API درواره پردازش هوش مصنوعی را انجام می‌دهد.
  • نتیجه در Queue جداگانه منتشر می‌شود.

RabbitMQ از Protocolهای مختلف پشتیبانی می‌کند، اما AMQP 0-9-1 یکی از مدل‌های اصلی و پرکاربرد آن است.

نسخه جاری RabbitMQ در زمان نگارش این مقاله 4.3.4 است. نسخه به‌روز و روش‌های نصب در مستندات رسمی نصب RabbitMQ اعلام می‌شود.

Message Queue چیست؟

Message Queue یا صف پیام واسطه‌ای میان تولیدکننده و مصرف‌کننده پیام است.

فرض کنید کاربر درخواست خلاصه‌سازی یک متن را ارسال می‌کند. در معماری Synchronous، FastAPI مستقیماً مدل را فراخوانی می‌کند:

کاربر
→ FastAPI
→ مدل هوش مصنوعی
→ FastAPI
→ کاربر

در معماری Queue-Based، FastAPI درخواست را به یک Job تبدیل می‌کند:

کاربر
→ FastAPI
→ RabbitMQ
→ Worker
→ مدل هوش مصنوعی

FastAPI می‌تواند بلافاصله پاسخ زیر را برگرداند:

{
  "job_id": "job-123",
  "status": "queued"
}

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

RabbitMQ چگونه کار می‌کند؟

برخلاف تصور اولیه، Producer در معماری استاندارد RabbitMQ معمولاً پیام را مستقیماً داخل Queue نمی‌نویسد. پیام ابتدا به Exchange فرستاده می‌شود.

معماری کامل‌تر:

Producer
→ Exchange
→ Binding
→ Queue
→ Consumer

Exchange با استفاده از نوع Exchange، Routing Key و Bindingها تعیین می‌کند پیام به کدام Queue یا Queueها ارسال شود.

مفاهیم اصلی RabbitMQ

Producer

Producer برنامه‌ای است که Message را منتشر می‌کند.

در پروژه ما، Endpoint مربوط به FastAPI پس از اعتبارسنجی درخواست، یک Message JSON تولید و به RabbitMQ ارسال می‌کند.

Consumer

Consumer برنامه‌ای است که Message را از Queue دریافت و پردازش می‌کند.

یک Queue می‌تواند چند Consumer داشته باشد. RabbitMQ پیام‌ها را میان Consumerهای فعال توزیع می‌کند.

Queue

Queue محلی است که Messageها تا زمان تحویل یا پردازش نگهداری می‌شوند.

نمونه نام Queueها:

ai.jobs.summarize
ai.jobs.completed
ai.jobs.dead
email.send
report.generate
image.process

Exchange

Exchange Message را از Producer دریافت و براساس Bindingها Route می‌کند.

Exchange می‌تواند Message را:

  • به یک Queue بفرستد.
  • به چند Queue بفرستد.
  • براساس Pattern انتخاب کند.
  • در صورت نبود Route مناسب برگرداند یا به مسیر دیگری هدایت کند.

Binding

Binding رابطه میان Exchange و Queue است.

برای مثال:

Exchange: ai.jobs
Routing Key: summarize
Queue: ai.jobs.summarize

این Binding به RabbitMQ می‌گوید Messageهایی که با Routing Key برابر summarize به Exchange ai.jobs می‌رسند، به Queue ai.jobs.summarize فرستاده شوند.

Routing Key

Routing Key یک رشته است که Producer هنگام انتشار Message مشخص می‌کند. Exchange براساس این مقدار و Bindingها مسیر Message را تعیین می‌کند.

نمونه Routing Keyها:

summarize
rewrite
translate
image.generate
document.analyze
notification.email

Channel

Channel یک ارتباط منطقی داخل Connection است. ساخت Connection شبکه برای هر Message پرهزینه است؛ بنابراین Client معمولاً یک Connection پایدار ایجاد و یک یا چند Channel روی آن باز می‌کند.

Connection و Channel نباید برای هر Request از ابتدا ساخته شوند. در FastAPI بهتر است در Lifespan برنامه ایجاد و هنگام Shutdown بسته شوند.

Virtual Host

Virtual Host یا vhost محیطی منطقی برای جداسازی Exchangeها، Queueها، Bindingها و دسترسی‌هاست.

برای مثال می‌توان محیط‌ها را جدا کرد:

/development
/staging
/production

در پروژه محلی این مقاله از vhost پیش‌فرض / استفاده می‌کنیم.

Message

Message علاوه بر Body می‌تواند Metadata نیز داشته باشد:

  • message_id
  • correlation_id
  • content_type
  • content_encoding
  • timestamp
  • headers
  • delivery_mode
  • expiration
  • reply_to
  • type

نمونه Message:

{
  "event_id": "2a905c48-cc15-416a-84a7-a7c39517ea61",
  "event_type": "ai.job.requested",
  "event_version": 1,
  "occurred_at": "2026-08-06T12:00:00Z",
  "payload": {
    "job_id": "job-123",
    "operation": "summarize",
    "text": "متن موردنظر کاربر"
  }
}

انواع Exchange در RabbitMQ

RabbitMQ چهار Exchange اصلی دارد:

  • Direct
  • Topic
  • Fanout
  • Headers

براساس مستندات Exchangeهای RabbitMQ، هر نوع Exchange منطق متفاوتی برای Routing دارد.

Direct Exchange

Direct Exchange پیام را به Queueهایی می‌فرستد که Binding Key آن‌ها دقیقاً با Routing Key برابر باشد.

مثال:

Routing Key: summarize

Binding:

summarize → ai.jobs.summarize

این نوع Exchange برای Task Queue و مسیریابی دقیق مناسب است.

Topic Exchange

Topic Exchange از Pattern Matching استفاده می‌کند. بخش‌های Routing Key با نقطه جدا می‌شوند.

دو Wildcard مهم:

  • * دقیقاً یک بخش را تطبیق می‌دهد.
  • # صفر یا چند بخش را تطبیق می‌دهد.

نمونه Routing Key:

ai.text.summarize
ai.text.translate
ai.image.generate

Bindingهای نمونه:

ai.text.*     → همه عملیات متنی تک‌مرحله‌ای
ai.#          → تمام پیام‌های مرتبط با هوش مصنوعی
*.image.*     → عملیات دارای بخش image

Topic Exchange برای معماری‌هایی که دسته‌بندی سلسله‌مراتبی Event دارند مناسب است.

Fanout Exchange

Fanout Exchange Routing Key را نادیده می‌گیرد و یک نسخه از Message را به همه Queueهای متصل می‌فرستد.

مثال:

user.registered
├── email.queue
├── analytics.queue
└── crm.queue

این الگو برای Publish/Subscribe مناسب است.

Headers Exchange

Headers Exchange به‌جای Routing Key از Headerهای Message استفاده می‌کند. این نوع Exchange نسبت به Direct و Topic کمتر استفاده می‌شود، اما برای بعضی قواعد چندشرطی مناسب است.

Default Exchange

RabbitMQ یک Direct Exchange پیش‌فرض با نام خالی دارد. در این حالت Routing Key باید با نام Queue برابر باشد.

نمونه:

await channel.default_exchange.publish(
    message,
    routing_key="ai.jobs.summarize",
)

برای پروژه‌های واقعی استفاده از Exchange نام‌گذاری‌شده معمولاً ساختار روشن‌تر و قابل‌توسعه‌تری ایجاد می‌کند.

RabbitMQ چه تفاوتی با Redis دارد؟

Redis یک In-Memory Data Store چندمنظوره است که می‌تواند برای Cache، Session، Rate Limit، Pub/Sub، Stream و بعضی Queueها استفاده شود.

RabbitMQ یک Message Broker تخصصی است که روی Routing، Queue، Acknowledgement، Delivery و ارتباط میان Producer و Consumer تمرکز دارد.

ویژگیRabbitMQRedis
کاربرد اصلیMessage BrokerIn-Memory Data Store
RoutingDirect، Topic، Fanout و Headersمحدودتر
Acknowledgementقابلیت اصلیوابسته به ساختار انتخابی
Dead Letterپشتیبانی داردنیازمند طراحی
Queue Managementتخصصییکی از کاربردها
Cacheکاربرد اصلی نیستبسیار مناسب
Message Priorityپشتیبانی در Classic Queueنیازمند طراحی
Pub/Subدارددارد
Persistent Messagingداردوابسته به Persistence
Management UIPlugin رسمی داردابزارهای جداگانه

اگر پروژه از قبل Redis دارد و Queue ساده‌ای نیاز دارد، Redis ممکن است کافی باشد. اگر Routing، Ack، DLQ و رفتار دقیق Message Broker اهمیت دارد، RabbitMQ انتخاب تخصصی‌تری است.

RabbitMQ چه تفاوتی با Kafka دارد؟

RabbitMQ و Kafka در بعضی سناریوها هم‌پوشانی دارند، اما مدل اصلی آن‌ها متفاوت است.

معیارRabbitMQKafka
مدل اصلیMessage Broker و QueueDistributed Event Log
حذف پس از پردازشمعمولاً پس از Ackبراساس Retention
Replayکاربرد اصلی نیستقابلیت مهم
Routingبسیار انعطاف‌پذیرعمدتاً Topic و Partition
ترتیبوابسته به Queue و Consumerدر هر Partition
Task Queueبسیار مناسبممکن است پیچیده‌تر باشد
Event Streaming حجیمممکن استبسیار مناسب
چند Consumer مستقلبا Queueهای جدابا Consumer Group
Message Priorityدر Classic Queueقابلیت پایه نیست
پیچیدگی اولیهمتوسطبیشتر

برای Jobهای پس‌زمینه، ارسال ایمیل، پردازش فایل و Task Routing، RabbitMQ معمولاً انتخاب طبیعی‌تری است. برای Event Streaming، Replay تاریخی و جریان‌های حجیم، Kafka بیشتر استفاده می‌شود.

Acknowledgement یا Ack چیست؟

وقتی RabbitMQ یک Message را به Consumer تحویل می‌دهد، باید بداند پردازش موفق بوده است یا خیر.

در Manual Acknowledgement، Consumer بعد از پردازش موفق Ack می‌فرستد:

RabbitMQ → Consumer
Consumer پردازش می‌کند
Consumer → Ack
RabbitMQ Message را حذف می‌کند

اگر Consumer قبل از Ack قطع شود، RabbitMQ می‌تواند Message تأییدنشده را دوباره در دسترس Consumer دیگری قرار دهد.

در Automatic Acknowledgement، Message بلافاصله پس از تحویل موفق تلقی می‌شود. اگر Consumer بعد از دریافت و پیش از پایان پردازش متوقف شود، امکان ازدست‌رفتن Job وجود دارد.

برای پردازش‌های مهم، Manual Ack انتخاب مناسب‌تری است.

راهنمای رسمی Ack و Publisher Confirm در صفحه Consumer Acknowledgements and Publisher Confirms قرار دارد.

تفاوت Ack، Reject و Nack

Ack

پردازش موفق بوده است:

basic.ack

Reject

پردازش یک Message شکست خورده است:

basic.reject

Consumer مشخص می‌کند Message دوباره وارد Queue شود یا خیر.

Nack

Nack مشابه Reject است، اما قابلیت‌های بیشتری مانند ردکردن چند Message را در بعضی Clientها فراهم می‌کند:

basic.nack

در هر دو حالت باید درباره requeue تصمیم بگیرید.

اگر requeue=true باشد، Message ممکن است دوباره وارد همان Queue شود. بدون کنترل Retry، این رفتار می‌تواند یک Loop بی‌پایان ایجاد کند.

اگر requeue=false باشد، Message در صورت وجود Dead Letter Exchange به آن منتقل می‌شود؛ در غیر این صورت ممکن است حذف شود.

Publisher Confirm چیست؟

Consumer Ack تأیید می‌کند که Consumer پیام را پردازش کرده است. Publisher Confirm مسئله متفاوتی را حل می‌کند: آیا Broker پیام منتشرشده را پذیرفته است؟

مسیر Publisher Confirm:

Producer → RabbitMQ
RabbitMQ → Confirm

بدون Publisher Confirm، Producer نمی‌تواند فقط براساس موفقیت ارسال روی Socket مطمئن شود Message توسط Broker پذیرفته شده است.

در پروژه این مقاله Publisher Confirm را فعال می‌کنیم.

Durable Queue و Persistent Message

برای افزایش دوام Message دو تنظیم مکمل وجود دارد.

Durable Queue

Queue پس از Restart شدن RabbitMQ دوباره ایجاد می‌شود:

durable=True

Persistent Message

Message برای نگهداری پایدار علامت‌گذاری می‌شود:

delivery_mode=aio_pika.DeliveryMode.PERSISTENT

Durable بودن Queue به‌تنهایی Message را Persistent نمی‌کند. Persistent بودن Message نیز اگر Queue موقتی باشد، ساختار Queue را پس از Restart حفظ نمی‌کند.

این تنظیم‌ها احتمال ازدست‌رفتن Message در Restart را کاهش می‌دهند، اما به‌تنهایی تضمین مطلق سراسری ایجاد نمی‌کنند. برای اطمینان بیشتر باید Publisher Confirm، نوع Queue، Replication، Storage و معماری Idempotent نیز بررسی شوند.

Prefetch چیست؟

Prefetch مشخص می‌کند RabbitMQ حداکثر چند Message تأییدنشده را هم‌زمان به Consumer تحویل دهد.

مثال:

await channel.set_qos(prefetch_count=5)

در این حالت Consumer حداکثر پنج Message پردازش‌نشده در اختیار دارد.

اگر Prefetch بسیار بزرگ باشد:

  • یک Worker ممکن است تعداد زیادی Message را در اختیار بگیرد.
  • توزیع بار نامتوازن می‌شود.
  • مصرف حافظه افزایش پیدا می‌کند.
  • هنگام توقف Worker تعداد زیادی Message باید دوباره تحویل داده شود.

اگر Prefetch بسیار کوچک باشد، ممکن است Throughput کاهش پیدا کند.

براساس مستندات Consumer Prefetch، مقدار مناسب باید براساس زمان پردازش، تعداد Workerها و منابع سیستم انتخاب شود.

برای Jobهای سنگین هوش مصنوعی، مقدار اولیه ۱ تا ۱۰ قابل آزمایش است؛ اما مقدار نهایی باید با Load Test تعیین شود.

Dead Letter Exchange چیست؟

Dead Letter Exchange یا DLX مسیری برای Messageهایی است که در Queue اصلی قابل پردازش نبوده‌اند.

Message ممکن است در این شرایط Dead Letter شود:

  • Consumer آن را با requeue=false رد کند.
  • Message منقضی شود.
  • Queue از محدودیت طول عبور کند.
  • در بعضی نوع Queueها Delivery Limit رد شود.

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

Producer
→ ai.jobs Exchange
→ ai.jobs.summarize Queue
→ Worker

اگر پردازش ناموفق بود:
ai.jobs.summarize
→ ai.jobs.dlx
→ ai.jobs.dead Queue

DLQ یا Dead Letter Queue محل بررسی، گزارش و پردازش کنترل‌شده Messageهای ناموفق است.

جزئیات رفتار DLX در مستندات Dead Letter Exchanges توضیح داده شده است.

نصب RabbitMQ با Docker

برای اجرای محلی RabbitMQ یک فایل docker-compose.yml بسازید:

services:
  rabbitmq:
    image: rabbitmq:4.3.4-management
    container_name: darvareh-rabbitmq
    hostname: rabbitmq
    environment:
      RABBITMQ_DEFAULT_USER: app
      RABBITMQ_DEFAULT_PASS: change-this-local-password
    ports:
      - "127.0.0.1:5672:5672"
      - "127.0.0.1:15672:15672"
    volumes:
      - rabbitmq_data:/var/lib/rabbitmq
    healthcheck:
      test:
        [
          "CMD",
          "rabbitmq-diagnostics",
          "-q",
          "ping"
        ]
      interval: 10s
      timeout: 5s
      retries: 10

volumes:
  rabbitmq_data:

RabbitMQ را اجرا کنید:

docker compose up -d

وضعیت Container:

docker compose ps

مشاهده Logها:

docker compose logs -f rabbitmq

بررسی سلامت:

docker exec darvareh-rabbitmq \
  rabbitmq-diagnostics -q ping

اگر پاسخ Ping succeeded دریافت کردید، RabbitMQ آماده است.

ورود به RabbitMQ Management UI

نسخه management شامل رابط مدیریتی RabbitMQ است.

آدرس:

http://127.0.0.1:15672

اطلاعات ورود محیط محلی:

Username: app
Password: change-this-local-password

در رابط مدیریتی می‌توانید این موارد را ببینید:

  • Connectionها
  • Channelها
  • Exchangeها
  • Queueها
  • Bindingها
  • Message Rate
  • Consumerها
  • تعداد Messageهای Ready
  • تعداد Messageهای Unacknowledged
  • Nodeها

این تنظیم فقط برای محیط محلی است. Management UI نباید بدون کنترل دسترسی مناسب در اینترنت عمومی منتشر شود.

ساخت پروژه عملی RabbitMQ با FastAPI و Python

در پروژه زیر یک Pipeline پردازش متن می‌سازیم:

FastAPI
→ Direct Exchange
→ Request Queue
→ AI Worker
→ API درواره
→ Completed Exchange
→ Completed Queue

اگر Worker نتواند Job را پس از Retry محدود پردازش کند:

Request Queue
→ Dead Letter Exchange
→ Dead Letter Queue

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

rabbitmq-ai-jobs/
├── docker-compose.yml
├── requirements.txt
├── .env
├── messaging.py
├── api.py
├── worker.py
└── result_consumer.py

نصب کتابخانه‌ها

فایل requirements.txt:

fastapi>=0.115,<1
uvicorn[standard]>=0.34,<1
aio-pika>=9.5,<10
openai>=1.100,<2
python-dotenv>=1.0,<2
pydantic>=2.10,<3

ساخت Virtual Environment:

python -m venv .venv

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

source .venv/bin/activate

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

.venv\Scripts\Activate.ps1

نصب وابستگی‌ها:

pip install -r requirements.txt

در این پروژه از aio-pika استفاده می‌کنیم. این کتابخانه Client غیرهم‌زمان RabbitMQ برای Python است و از connect_robust، بازیابی Connection و Publisher Confirm پشتیبانی می‌کند. راهنمای آن در مستندات aio-pika قرار دارد.

تنظیم متغیرهای محیطی

فایل .env:

RABBITMQ_URL=amqp://app:change-this-local-password@127.0.0.1:5672/

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

فایل .env را در Git ثبت نکنید.

نمونه .gitignore:

.env
.venv/
__pycache__/
*.pyc

تعریف نام Exchange و Queueها

فایل messaging.py:

import aio_pika
from aio_pika.abc import (
    AbstractChannel,
    AbstractExchange,
    AbstractQueue,
)

REQUEST_EXCHANGE = "ai.jobs"
REQUEST_QUEUE = "ai.jobs.summarize"
REQUEST_ROUTING_KEY = "summarize"

COMPLETED_EXCHANGE = "ai.jobs.completed"
COMPLETED_QUEUE = "ai.jobs.completed"
COMPLETED_ROUTING_KEY = "completed"

DEAD_LETTER_EXCHANGE = "ai.jobs.dlx"
DEAD_LETTER_QUEUE = "ai.jobs.dead"
DEAD_LETTER_ROUTING_KEY = "failed"


async def declare_topology(
    channel: AbstractChannel,
) -> dict[str, AbstractExchange | AbstractQueue]:
    dead_exchange = await channel.declare_exchange(
        DEAD_LETTER_EXCHANGE,
        aio_pika.ExchangeType.DIRECT,
        durable=True,
    )

    dead_queue = await channel.declare_queue(
        DEAD_LETTER_QUEUE,
        durable=True,
    )

    await dead_queue.bind(
        dead_exchange,
        routing_key=DEAD_LETTER_ROUTING_KEY,
    )

    request_exchange = await channel.declare_exchange(
        REQUEST_EXCHANGE,
        aio_pika.ExchangeType.DIRECT,
        durable=True,
    )

    request_queue = await channel.declare_queue(
        REQUEST_QUEUE,
        durable=True,
        arguments={
            "x-dead-letter-exchange": (
                DEAD_LETTER_EXCHANGE
            ),
            "x-dead-letter-routing-key": (
                DEAD_LETTER_ROUTING_KEY
            ),
        },
    )

    await request_queue.bind(
        request_exchange,
        routing_key=REQUEST_ROUTING_KEY,
    )

    completed_exchange = await channel.declare_exchange(
        COMPLETED_EXCHANGE,
        aio_pika.ExchangeType.DIRECT,
        durable=True,
    )

    completed_queue = await channel.declare_queue(
        COMPLETED_QUEUE,
        durable=True,
    )

    await completed_queue.bind(
        completed_exchange,
        routing_key=COMPLETED_ROUTING_KEY,
    )

    return {
        "request_exchange": request_exchange,
        "request_queue": request_queue,
        "completed_exchange": completed_exchange,
        "completed_queue": completed_queue,
        "dead_exchange": dead_exchange,
        "dead_queue": dead_queue,
    }

تعریف Topology در یک تابع مشترک باعث می‌شود API و Worker از نام‌ها و تنظیمات یکسان استفاده کنند.

در پروژه‌های بزرگ‌تر بهتر است Queue، Exchange و Policyها با Infrastructure as Code یا فرایند مدیریت‌شده ایجاد شوند.

ساخت Producer با FastAPI

فایل api.py:

import json
import os
from contextlib import asynccontextmanager
from datetime import datetime, timezone
from typing import Literal
from uuid import uuid4

import aio_pika
from aio_pika.abc import (
    AbstractChannel,
    AbstractRobustConnection,
)
from dotenv import load_dotenv
from fastapi import FastAPI, HTTPException, Request
from pydantic import BaseModel, Field

from messaging import (
    REQUEST_ROUTING_KEY,
    declare_topology,
)

load_dotenv()

RABBITMQ_URL = os.getenv(
    "RABBITMQ_URL",
    "amqp://app:change-this-local-password"
    "@127.0.0.1:5672/",
)


class AIJobRequest(BaseModel):
    text: str = Field(
        min_length=2,
        max_length=50_000,
    )

    operation: Literal[
        "summarize",
        "rewrite",
        "classify",
    ] = "summarize"

    user_id: str | None = Field(
        default=None,
        max_length=100,
    )


def utc_now() -> str:
    return datetime.now(timezone.utc).isoformat()


@asynccontextmanager
async def lifespan(app: FastAPI):
    connection = await aio_pika.connect_robust(
        RABBITMQ_URL,
        client_properties={
            "connection_name": "ai-job-api",
        },
    )

    channel = await connection.channel(
        publisher_confirms=True,
        on_return_raises=True,
    )

    topology = await declare_topology(channel)

    app.state.rabbitmq_connection = connection
    app.state.rabbitmq_channel = channel
    app.state.request_exchange = topology[
        "request_exchange"
    ]

    try:
        yield
    finally:
        await channel.close()
        await connection.close()


app = FastAPI(
    title="RabbitMQ AI Job API",
    version="1.0.0",
    lifespan=lifespan,
)


@app.get("/health")
async def health(request: Request):
    connection: AbstractRobustConnection = (
        request.app.state.rabbitmq_connection
    )

    channel: AbstractChannel = (
        request.app.state.rabbitmq_channel
    )

    return {
        "status": "ok",
        "rabbitmq_connection_closed": (
            connection.is_closed
        ),
        "rabbitmq_channel_closed": channel.is_closed,
    }


@app.post("/jobs", status_code=202)
async def create_job(
    payload: AIJobRequest,
    request: Request,
):
    job_id = str(uuid4())
    event_id = str(uuid4())

    event = {
        "event_id": event_id,
        "event_type": "ai.job.requested",
        "event_version": 1,
        "occurred_at": utc_now(),
        "trace_id": job_id,
        "payload": {
            "job_id": job_id,
            "user_id": payload.user_id,
            "operation": payload.operation,
            "text": payload.text,
        },
    }

    message = aio_pika.Message(
        body=json.dumps(
            event,
            ensure_ascii=False,
        ).encode("utf-8"),
        content_type="application/json",
        content_encoding="utf-8",
        delivery_mode=(
            aio_pika.DeliveryMode.PERSISTENT
        ),
        message_id=event_id,
        correlation_id=job_id,
        timestamp=datetime.now(timezone.utc),
        type="ai.job.requested",
        headers={
            "event_version": 1,
            "operation": payload.operation,
        },
    )

    exchange = request.app.state.request_exchange

    try:
        confirmed = await exchange.publish(
            message,
            routing_key=REQUEST_ROUTING_KEY,
            mandatory=True,
        )

        if confirmed is False:
            raise RuntimeError(
                "Message was not confirmed"
            )

    except Exception as exc:
        raise HTTPException(
            status_code=503,
            detail="Job could not be queued",
        ) from exc

    return {
        "job_id": job_id,
        "event_id": event_id,
        "status": "queued",
    }

نکات مهم این Producer:

  • Connection فقط یک‌بار در Lifespan ساخته می‌شود.
  • از connect_robust برای Reconnect استفاده شده است.
  • Publisher Confirm فعال است.
  • Message با Delivery Mode پایدار منتشر می‌شود.
  • mandatory=True باعث می‌شود Message بدون Route به‌سادگی نادیده گرفته نشود.
  • پاسخ 202 Accepted نشان می‌دهد Job پذیرفته شده، اما هنوز تکمیل نشده است.
  • برای هر Job شناسه job_id و برای هر Message شناسه event_id داریم.

ساخت Worker هوش مصنوعی

فایل worker.py:

import asyncio
import json
import os
from datetime import datetime, timezone
from uuid import uuid4

import aio_pika
from aio_pika.abc import AbstractIncomingMessage
from dotenv import load_dotenv
from openai import AsyncOpenAI

from messaging import (
    COMPLETED_ROUTING_KEY,
    REQUEST_QUEUE,
    declare_topology,
)

load_dotenv()

RABBITMQ_URL = os.getenv(
    "RABBITMQ_URL",
    "amqp://app:change-this-local-password"
    "@127.0.0.1:5672/",
)

DARVAREH_API_KEY = os.getenv(
    "DARVAREH_API_KEY",
    "",
)

DARVAREH_MODEL_ID = os.getenv(
    "DARVAREH_MODEL_ID",
    "YOUR_MODEL_ID",
)

DARVAREH_BASE_URL = os.getenv(
    "DARVAREH_BASE_URL",
    "https://api.darvareh.ir/v1",
)

if not DARVAREH_API_KEY:
    raise RuntimeError(
        "DARVAREH_API_KEY is not configured"
    )

ai_client = AsyncOpenAI(
    api_key=DARVAREH_API_KEY,
    base_url=DARVAREH_BASE_URL,
    timeout=60,
    max_retries=2,
)


def utc_now() -> str:
    return datetime.now(timezone.utc).isoformat()


def build_messages(
    operation: str,
    text: str,
) -> list[dict[str, str]]:
    instructions = {
        "summarize": (
            "متن را دقیق و کوتاه خلاصه کن. "
            "هیچ اطلاعاتی به متن اضافه نکن."
        ),
        "rewrite": (
            "متن را به فارسی روان و حرفه‌ای "
            "بازنویسی کن، بدون تغییر مفهوم."
        ),
        "classify": (
            "موضوع اصلی متن را مشخص کن. "
            "خروجی را فقط در یک عبارت کوتاه بنویس."
        ),
    }

    instruction = instructions.get(operation)

    if instruction is None:
        raise ValueError(
            f"Unsupported operation: {operation}"
        )

    return [
        {
            "role": "system",
            "content": instruction,
        },
        {
            "role": "user",
            "content": text,
        },
    ]


async def call_ai_with_retry(
    operation: str,
    text: str,
    max_attempts: int = 3,
) -> str:
    last_error: Exception | None = None

    for attempt in range(1, max_attempts + 1):
        try:
            response = (
                await ai_client.chat.completions.create(
                    model=DARVAREH_MODEL_ID,
                    temperature=0.2,
                    max_tokens=800,
                    messages=build_messages(
                        operation=operation,
                        text=text,
                    ),
                )
            )

            result = response.choices[
                0
            ].message.content

            if not result:
                raise ValueError(
                    "Model returned empty content"
                )

            return result

        except Exception as exc:
            last_error = exc

            if attempt < max_attempts:
                delay = 2 ** (attempt - 1)
                await asyncio.sleep(delay)

    raise RuntimeError(
        "AI processing failed after retries"
    ) from last_error


async def process_message(
    message: AbstractIncomingMessage,
    completed_exchange,
) -> None:
    async with message.process(
        requeue=False,
        ignore_processed=False,
    ):
        event = json.loads(
            message.body.decode("utf-8")
        )

        payload = event["payload"]

        job_id = payload["job_id"]
        operation = payload["operation"]
        text = payload["text"]

        result = await call_ai_with_retry(
            operation=operation,
            text=text,
        )

        completed_event = {
            "event_id": str(uuid4()),
            "event_type": "ai.job.completed",
            "event_version": 1,
            "occurred_at": utc_now(),
            "trace_id": event["trace_id"],
            "causation_id": event["event_id"],
            "payload": {
                "job_id": job_id,
                "operation": operation,
                "result": result,
                "model": DARVAREH_MODEL_ID,
            },
        }

        completed_message = aio_pika.Message(
            body=json.dumps(
                completed_event,
                ensure_ascii=False,
            ).encode("utf-8"),
            content_type="application/json",
            content_encoding="utf-8",
            delivery_mode=(
                aio_pika.DeliveryMode.PERSISTENT
            ),
            message_id=completed_event["event_id"],
            correlation_id=job_id,
            timestamp=datetime.now(timezone.utc),
            type="ai.job.completed",
        )

        confirmed = await completed_exchange.publish(
            completed_message,
            routing_key=COMPLETED_ROUTING_KEY,
            mandatory=True,
        )

        if confirmed is False:
            raise RuntimeError(
                "Completed event was not confirmed"
            )

        print(
            f"Job completed: {job_id}"
        )


async def main() -> None:
    connection = await aio_pika.connect_robust(
        RABBITMQ_URL,
        client_properties={
            "connection_name": "ai-job-worker",
        },
    )

    channel = await connection.channel(
        publisher_confirms=True,
        on_return_raises=True,
    )

    await channel.set_qos(
        prefetch_count=5,
    )

    topology = await declare_topology(channel)

    request_queue = topology["request_queue"]
    completed_exchange = topology[
        "completed_exchange"
    ]

    print(
        f"Worker is consuming {REQUEST_QUEUE}"
    )

    try:
        async with request_queue.iterator() as iterator:
            async for message in iterator:
                try:
                    await process_message(
                        message,
                        completed_exchange,
                    )

                except asyncio.CancelledError:
                    raise

                except Exception as exc:
                    print(
                        "Message failed and was "
                        f"dead-lettered: {exc}"
                    )

    finally:
        await channel.close()
        await connection.close()
        await ai_client.close()


if __name__ == "__main__":
    asyncio.run(main())

در این Worker:

  1. Message با Manual Ack دریافت می‌شود.
  2. متن حداکثر سه بار پردازش می‌شود.
  3. نتیجه در Exchange خروجی منتشر می‌شود.
  4. Publisher Confirm نتیجه بررسی می‌شود.
  5. فقط پس از خروج موفق از Context Manager، Message اصلی Ack می‌شود.
  6. اگر پردازش شکست بخورد، Message با requeue=false رد و به DLX منتقل می‌شود.

این طراحی جلوی Retry بی‌نهایت در Queue اصلی را می‌گیرد.

بااین‌حال اگر Worker بعد از انتشار نتیجه و پیش از Ack متوقف شود، Message اصلی ممکن است دوباره تحویل داده شود. بنابراین Consumer نتیجه و Storage نهایی باید براساس event_id یا job_id رفتار Idempotent داشته باشند.

ساخت Consumer نتایج

فایل result_consumer.py:

import asyncio
import json
import os

import aio_pika
from dotenv import load_dotenv

from messaging import declare_topology

load_dotenv()

RABBITMQ_URL = os.getenv(
    "RABBITMQ_URL",
    "amqp://app:change-this-local-password"
    "@127.0.0.1:5672/",
)


async def main() -> None:
    connection = await aio_pika.connect_robust(
        RABBITMQ_URL,
        client_properties={
            "connection_name": "ai-result-consumer",
        },
    )

    channel = await connection.channel()

    await channel.set_qos(
        prefetch_count=10,
    )

    topology = await declare_topology(channel)
    completed_queue = topology["completed_queue"]

    print("Waiting for completed AI jobs")

    try:
        async with completed_queue.iterator() as iterator:
            async for message in iterator:
                async with message.process(
                    requeue=False,
                ):
                    event = json.loads(
                        message.body.decode("utf-8")
                    )

                    payload = event["payload"]

                    print(
                        "\nJob ID:",
                        payload["job_id"],
                    )

                    print(
                        "Operation:",
                        payload["operation"],
                    )

                    print(
                        "Model:",
                        payload["model"],
                    )

                    print(
                        "Result:",
                        payload["result"],
                    )

    finally:
        await channel.close()
        await connection.close()


if __name__ == "__main__":
    asyncio.run(main())

در یک پروژه واقعی، Consumer نتیجه معمولاً خروجی را در PostgreSQL یا Storage دیگری ذخیره می‌کند تا API بتواند وضعیت Job را برگرداند:

GET /jobs/{job_id}

پاسخ نمونه:

{
  "job_id": "job-123",
  "status": "completed",
  "result": "خلاصه تولیدشده"
}

اجرای کامل پروژه

ابتدا RabbitMQ را اجرا کنید:

docker compose up -d

FastAPI را اجرا کنید:

uvicorn api:app --reload

Worker را در Terminal دوم اجرا کنید:

python worker.py

Consumer نتیجه را در Terminal سوم اجرا کنید:

python result_consumer.py

اکنون یک Job ارسال کنید:

curl -X POST "http://127.0.0.1:8000/jobs" \
  -H "Content-Type: application/json" \
  -d '{
    "text": "هوش مصنوعی می‌تواند فرایندهای تکراری را سریع‌تر کند، اما کیفیت خروجی به انتخاب مدل، طراحی پرامپت، داده ورودی و ارزیابی مستمر وابسته است.",
    "operation": "summarize",
    "user_id": "user-123"
  }'

پاسخ API:

{
  "job_id": "fe241bb8-17d1-4b51-94d2-351073f37dc1",
  "event_id": "a320314e-736a-4d0f-a5bb-9f638f380abd",
  "status": "queued"
}

پس از چند لحظه، خروجی در result_consumer.py نمایش داده می‌شود.

برای دریافت کلید API و استفاده از مدل‌های مختلف، در درواره ثبت‌نام کنید. شناسه مدل‌ها، قابلیت‌ها و قیمت به‌روز آن‌ها در صفحه مدل‌های درواره قرار دارد.

مشاهده Queueها در Management UI

وارد این آدرس شوید:

http://127.0.0.1:15672

در بخش Queues باید این Queueها را ببینید:

ai.jobs.summarize
ai.jobs.completed
ai.jobs.dead

ستون‌های مهم:

  • Ready: Messageهایی که منتظر Consumer هستند.
  • Unacked: Messageهایی که تحویل داده شده‌اند اما هنوز Ack نشده‌اند.
  • Total: مجموع Ready و Unacked.
  • Consumers: تعداد Consumerهای متصل.
  • Message rates: نرخ ورود، تحویل و Ack.

اگر Worker را متوقف و چند Job ارسال کنید، مقدار Ready در Queue درخواست افزایش پیدا می‌کند. بعد از اجرای Worker، Messageها پردازش می‌شوند.

آزمایش Dead Letter Queue

برای آزمایش DLQ می‌توانید موقتاً مقدار زیر را در .env قرار دهید:

DARVAREH_MODEL_ID=INVALID_MODEL_ID

Worker را Restart و یک Job ارسال کنید.

Worker پس از Retry محدود، Message را رد می‌کند. Message باید از Queue اصلی به Queue زیر منتقل شود:

ai.jobs.dead

پس از اصلاح مشکل، Messageهای DLQ را بدون بررسی مستقیم به Queue اصلی برنگردانید. ابتدا علت خطا را مشخص کنید:

  • آیا Payload نامعتبر بوده است؟
  • آیا مدل اشتباه تنظیم شده است؟
  • آیا خطا موقت یا دائمی است؟
  • آیا پردازش مجدد Side Effect تکراری ایجاد می‌کند؟
  • آیا Schema تغییر کرده است؟
  • آیا Message هنوز از نظر زمانی معتبر است؟

Retry استاندارد چگونه طراحی می‌شود؟

در پروژه آموزشی، Retry داخل Worker انجام می‌شود:

Attempt 1
→ 1 second delay
Attempt 2
→ 2 seconds delay
Attempt 3
→ DLQ

برای سیستم بزرگ‌تر می‌توان از Retry Queueهای زمان‌دار استفاده کرد:

ai.jobs.retry.10s
ai.jobs.retry.1m
ai.jobs.retry.10m
ai.jobs.dead

با TTL و Dead Letter Exchange می‌توان Message را پس از تأخیر به Queue اصلی بازگرداند.

ساختار مفهومی:

Main Queue
→ Retry Exchange
→ Retry Queue دارای TTL
→ پس از Expire
→ Main Exchange
→ Main Queue

در تنظیمات Production بهتر است Policyهای RabbitMQ برای DLX و TTL بررسی شوند تا تغییر تنظیمات بدون Redeploy همه برنامه‌ها امکان‌پذیر باشد.

Retry باید محدود باشد. Retry نامحدود می‌تواند:

  • Queue را پر کند.
  • هزینه API هوش مصنوعی را افزایش دهد.
  • سرویس خارجی را تحت فشار قرار دهد.
  • Messageهای سالم را پشت Message خراب نگه دارد.
  • تشخیص خطای دائمی را دشوار کند.

Idempotency در RabbitMQ

RabbitMQ می‌تواند در بعضی شرایط Message را دوباره تحویل دهد. برای مثال:

  1. Worker Message را پردازش می‌کند.
  2. نتیجه را در Database ذخیره می‌کند.
  3. پیش از Ack متوقف می‌شود.
  4. RabbitMQ Message را به Worker دیگری می‌دهد.
  5. عملیات دوباره اجرا می‌شود.

بنابراین هر Job باید شناسه یکتا داشته باشد:

job_id
event_id

در PostgreSQL می‌توان Constraint یکتا تعریف کرد:

CREATE TABLE processed_events (
    event_id UUID PRIMARY KEY,
    processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

یا نتیجه Job را با شناسه یکتا ذخیره کرد:

CREATE TABLE ai_jobs (
    job_id UUID PRIMARY KEY,
    event_id UUID UNIQUE NOT NULL,
    status TEXT NOT NULL,
    operation TEXT NOT NULL,
    result TEXT,
    error_message TEXT,
    created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    completed_at TIMESTAMPTZ
);

الگوی پردازش:

دریافت Message
→ بررسی event_id
→ اگر قبلاً پردازش شده: Ack
→ اگر جدید است: پردازش
→ ذخیره نتیجه
→ ثبت event_id
→ Ack

اگر ثبت نتیجه و ثبت Idempotency جداگانه انجام شوند، ممکن است میان آن‌ها خطا رخ دهد. بهتر است تا حد امکان هر دو در یک Transaction پایگاه داده انجام شوند.

Outbox Pattern

یکی از چالش‌های مهم Producer این است:

  1. داده در Database ذخیره شود.
  2. Message در RabbitMQ منتشر شود.

اگر Database موفق و RabbitMQ ناموفق باشد، وضعیت ناقص ایجاد می‌شود.

Outbox Pattern این مسئله را کاهش می‌دهد:

Transaction پایگاه داده:
├── ثبت Job
└── ثبت Outbox Event

یک Outbox Worker رکوردهای منتشرنشده را می‌خواند، آن‌ها را در RabbitMQ منتشر و پس از دریافت Publisher Confirm وضعیت را به Published تغییر می‌دهد.

نمونه جدول:

CREATE TABLE outbox_events (
    id UUID PRIMARY KEY,
    event_type TEXT NOT NULL,
    payload JSONB NOT NULL,
    status TEXT NOT NULL DEFAULT 'pending',
    attempts INTEGER NOT NULL DEFAULT 0,
    created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    published_at TIMESTAMPTZ
);

Outbox Pattern برای عملیات مهم، قابل‌اعتمادتر از انتشار مستقیم Message پس از Commit جداگانه Database است.

Classic Queue و Quorum Queue

RabbitMQ چند نوع Queue ارائه می‌کند. دو نوع مهم عبارت‌اند از:

  • Classic Queue
  • Quorum Queue

Classic Queue

برای بسیاری از سناریوهای عمومی، Queueهای موقت، Priority Queue و بعضی قابلیت‌های سنتی استفاده می‌شود.

Quorum Queue

Quorum Queue برای Queueهای Replicated و پردازش‌های مهم طراحی شده است و از الگوریتم اجماع برای Replication استفاده می‌کند.

در Cluster Production، Quorum Queue برای Jobهای مهم قابل بررسی است. براساس مستندات Quorum Queues، این Queueها رفتار و هزینه منابع متفاوتی دارند و باید متناسب با تعداد Nodeها و نیاز دوام انتخاب شوند.

برای تعریف Quorum Queue:

queue = await channel.declare_queue(
    "ai.jobs.summarize",
    durable=True,
    arguments={
        "x-queue-type": "quorum",
    },
)

نوع Queue را بدون بررسی سازگاری قابلیت‌ها تغییر ندهید. Classic Queue و Quorum Queue از نظر Replication، Priority، DLX، محدودیت‌ها و مصرف منابع رفتار یکسانی ندارند.

مقیاس‌کردن Workerها

برای افزایش ظرفیت می‌توانید چند نمونه از Worker را اجرا کنید:

python worker.py
python worker.py
python worker.py

همه Workerها از Queue یکسان مصرف می‌کنند. RabbitMQ Messageها را میان آن‌ها توزیع می‌کند.

اگر prefetch_count=5 باشد، هر Worker حداکثر پنج Message تأییدنشده دریافت می‌کند.

ظرفیت تقریبی Workerها به این عوامل وابسته است:

  • زمان پاسخ مدل
  • Rate Limit
  • تعداد Token
  • نوع عملیات
  • تعداد Connection و Channel
  • مقدار Prefetch
  • Concurrency داخلی
  • منابع CPU و Memory
  • محدودیت هزینه
  • تعداد Messageهای ورودی

افزایش Worker بدون توجه به Rate Limit API می‌تواند باعث افزایش خطای 429، Retry و هزینه عملیاتی شود.

Backpressure

اگر Producer سریع‌تر از Consumer Message تولید کند، طول Queue افزایش پیدا می‌کند. این وضعیت Backpressure را نشان می‌دهد.

شاخص‌های مهم:

  • Queue Depth
  • Message Publish Rate
  • Message Delivery Rate
  • Ack Rate
  • تعداد Unacked
  • عمر قدیمی‌ترین Message
  • تعداد Consumer
  • زمان متوسط پردازش
  • Retry Rate
  • DLQ Rate

راهکارها:

  • افزایش Worker به‌صورت کنترل‌شده
  • استفاده از مدل سریع‌تر
  • محدودکردن طول ورودی
  • Batch کردن عملیات سازگار
  • Route کردن Jobهای ساده و پیچیده
  • اعمال Rate Limit در Producer
  • اعلام زمان تقریبی انتظار
  • تعیین محدودیت طول Queue
  • اولویت‌بندی Jobها
  • بررسی Bottleneck واقعی پیش از Scale

انتخاب مدل مناسب برای Worker

همه Jobها نباید با یک مدل پردازش شوند.

نمونه Routing منطقی:

عملیاتنوع مدل پیشنهادی
دسته‌بندی سادهمدل سریع و اقتصادی
خلاصه‌سازی کوتاهمدل سبک متنی
تحلیل پیچیدهمدل Reasoning
پردازش تصویرمدل Vision
تولید کدمدل Coding
متن طولانیمدل دارای Context مناسب

شناسه مدل را از Configuration معتبر بخوانید. اجازه ندهید کاربر هر رشته دلخواهی را مستقیماً به پارامتر model ارسال کند.

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

مانیتورینگ RabbitMQ

برای Production این شاخص‌ها را مانیتور کنید:

شاخص‌های Queue

  • Ready Messages
  • Unacknowledged Messages
  • Total Messages
  • Publish Rate
  • Delivery Rate
  • Ack Rate
  • Redelivery Rate
  • Consumer Count
  • طول عمر Message
  • تعداد Messageهای DLQ

شاخص‌های Node

  • Memory Usage
  • Disk Usage
  • File Descriptor
  • Socket Count
  • Connection Count
  • Channel Count
  • Network Throughput
  • Erlang Process Count
  • Alarmهای Memory و Disk

شاخص‌های Worker

  • تعداد Job موفق
  • تعداد Job ناموفق
  • زمان پردازش
  • تعداد Retry
  • تعداد Timeout
  • نرخ Redelivery
  • تعداد Messageهای تکراری
  • زمان پاسخ مدل
  • تعداد Token
  • هزینه هر Job
  • P50، P95 و P99 Latency

لاگ‌گذاری ساختاریافته

برای هر Job این فیلدها را ثبت کنید:

event_id
job_id
trace_id
message_id
correlation_id
routing_key
exchange
queue
operation
model
attempt
duration_ms
status
error_type

نمونه Log:

{
  "level": "info",
  "event": "ai_job_completed",
  "job_id": "job-123",
  "event_id": "event-456",
  "queue": "ai.jobs.summarize",
  "operation": "summarize",
  "model": "YOUR_MODEL_ID",
  "duration_ms": 2450,
  "status": "completed"
}

کلید API، Credentialهای RabbitMQ و محتوای حساس کاربر را وارد Log نکنید.

تست RabbitMQ Pipeline

Unit Test

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

  • Validation ورودی
  • ساخت Event
  • انتخاب Routing Key
  • ساخت Prompt
  • Retry Policy
  • پردازش پاسخ مدل
  • مدیریت خطا
  • Idempotency

Integration Test

RabbitMQ آزمایشی اجرا و مسیر کامل را بررسی کنید:

HTTP Request
→ Exchange
→ Request Queue
→ Worker
→ API Mock
→ Completed Exchange
→ Completed Queue

تست خطا

این سناریوها را آزمایش کنید:

  • RabbitMQ هنگام Publish در دسترس نیست.
  • Queue Binding حذف شده است.
  • Message JSON نامعتبر است.
  • Worker هنگام پردازش متوقف می‌شود.
  • API مدل Timeout می‌دهد.
  • مدل خروجی خالی برمی‌گرداند.
  • Publisher نتیجه Confirm دریافت نمی‌کند.
  • Message وارد DLQ می‌شود.
  • Message تکراری دوباره تحویل داده می‌شود.

Load Test

اندازه‌گیری کنید:

  • چند Job در ثانیه Publish می‌شود؟
  • چند Job در ثانیه پردازش می‌شود؟
  • Queue با چه سرعتی رشد می‌کند؟
  • چند Worker برای بار عادی لازم است؟
  • زمان انتظار P95 چقدر است؟
  • چه زمانی Rate Limit مدل فعال می‌شود؟
  • بعد از قطعی، Queue در چه مدت تخلیه می‌شود؟

خطاهای رایج RabbitMQ

خطای Connection Refused

موارد زیر را بررسی کنید:

docker compose ps
docker compose logs rabbitmq
docker exec darvareh-rabbitmq \
  rabbitmq-diagnostics -q ping

اگر برنامه Python روی Host اجرا می‌شود:

127.0.0.1:5672

اگر برنامه داخل Docker Compose اجرا می‌شود:

rabbitmq:5672

داخل Container، localhost به همان Container اشاره می‌کند، نه Container مربوط به RabbitMQ.

خطای Authentication Failed

مقدار Username و Password در RABBITMQ_URL باید با Docker Compose یکسان باشد.

نمونه:

amqp://app:change-this-local-password@127.0.0.1:5672/

پس از ساخته‌شدن Volume، تغییر Environment Variable لزوماً User قبلی را تغییر نمی‌دهد. برای محیط محلی باید Credential را از طریق RabbitMQ مدیریت کنید یا Volume آزمایشی را با آگاهی از حذف داده بازسازی کنید.

Message منتشر می‌شود اما وارد Queue نمی‌شود

این موارد را بررسی کنید:

  • نام Exchange
  • نوع Exchange
  • Routing Key
  • Binding
  • نام Queue
  • فعال‌بودن mandatory
  • Durable بودن تعاریف سازگار
  • خطاهای Channel

در Management UI بخش Bindings را بررسی کنید.

تعداد Unacked زیاد می‌شود

دلایل احتمالی:

  • Consumer Ack نمی‌فرستد.
  • پردازش بسیار کند است.
  • Prefetch بیش‌ازحد بزرگ است.
  • Worker گیر کرده است.
  • درخواست مدل Timeout طولانی دارد.
  • Exception مدیریت نشده است.

Prefetch را کاهش دهید و زمان پردازش Worker را اندازه‌گیری کنید.

Message مرتب Requeue می‌شود

اگر Consumer همیشه requeue=true بفرستد، Message خراب ممکن است بی‌نهایت تکرار شود.

راهکار:

  • Retry محدود
  • تشخیص خطای موقت و دائمی
  • استفاده از Retry Queue
  • انتقال نهایی به DLQ
  • ثبت تعداد Attempt
  • Alert روی افزایش Redelivery

Message بعد از Restart باقی نمی‌ماند

بررسی کنید:

  • Queue با durable=True ساخته شده است.
  • Message با Delivery Mode پایدار منتشر شده است.
  • Exchange Durable است.
  • Queue موقتی یا Auto Delete نیست.
  • Storage RabbitMQ Volume دارد.
  • Publisher Confirm استفاده شده است.

خطای PRECONDITION_FAILED

این خطا معمولاً وقتی رخ می‌دهد که Exchange یا Queue موجود را با تنظیمات متفاوت دوباره Declare می‌کنید.

برای مثال Queue قبلاً Durable بوده و اکنون با durable=False Declare می‌شود.

نام یک Queue باید همیشه با تنظیمات سازگار Declare شود. برای تغییر نوع یا Argumentهای ناسازگار معمولاً باید Queue جدید با نام نسخه‌بندی‌شده ایجاد کنید.

Worker پیام دریافت نمی‌کند

موارد زیر را بررسی کنید:

  • Consumer به Queue درست متصل است.
  • Queue دارای Message Ready است.
  • Consumer دیگری Messageها را دریافت نمی‌کند.
  • Connection و Channel باز هستند.
  • Worker روی vhost درست قرار دارد.
  • Binding درست است.
  • Prefetch و Ack مشکل ندارند.

چک‌لیست Production

پیش از راه‌اندازی Production این موارد را بررسی کنید:

  • Exchangeها و Queueها نام‌گذاری مشخص دارند.
  • Topology نسخه‌بندی شده است.
  • Queueها در صورت نیاز Durable هستند.
  • Messageهای مهم Persistent هستند.
  • Publisher Confirm فعال است.
  • mandatory و پیام‌های Unroutable مدیریت می‌شوند.
  • Manual Ack استفاده می‌شود.
  • Ack فقط پس از Side Effect موفق ارسال می‌شود.
  • Prefetch با Load Test تنظیم شده است.
  • Retry محدود است.
  • Retry Loop بی‌نهایت وجود ندارد.
  • Dead Letter Exchange و DLQ ساخته شده‌اند.
  • DLQ مانیتور و Alert دارد.
  • Eventها event_id و job_id دارند.
  • Consumerها Idempotent هستند.
  • Outbox Pattern برای عملیات مهم بررسی شده است.
  • Timeout فراخوانی API مشخص است.
  • Shutdown برنامه Graceful است.
  • Connection برای هر Request ساخته نمی‌شود.
  • Queue Depth و Unacked مانیتور می‌شوند.
  • Messageهای بزرگ در Queue قرار نمی‌گیرند.
  • فایل‌ها در Object Storage ذخیره می‌شوند.
  • Credentialها در متغیر محیطی نگهداری می‌شوند.
  • مدل از فهرست معتبر انتخاب می‌شود.
  • هزینه هر Job اندازه‌گیری می‌شود.
  • برنامه بازیابی از قطعی آزمایش شده است.
  • Management UI در معرض عمومی قرار ندارد.

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

RabbitMQ چیست؟

RabbitMQ یک Message Broker است که Messageها را از Producer دریافت و براساس Exchange، Routing Key و Binding به Queueهای مناسب هدایت می‌کند تا Consumerها آن‌ها را پردازش کنند.

RabbitMQ برای چه کاری استفاده می‌شود؟

RabbitMQ برای Background Job، ارسال ایمیل، پردازش فایل، ارتباط Microserviceها، Task Queue، Notification، پردازش هوش مصنوعی و معماری Message-Driven استفاده می‌شود.

Queue در RabbitMQ چیست؟

Queue محلی است که Messageها تا زمان تحویل و تأیید پردازش توسط Consumer نگهداری می‌شوند.

Exchange چیست؟

Exchange Message را از Producer دریافت و براساس نوع Exchange، Routing Key و Binding به Queueهای مناسب Route می‌کند.

Routing Key چیست؟

Routing Key رشته‌ای است که Producer هنگام انتشار Message مشخص می‌کند. Exchange از آن برای انتخاب مسیر Message استفاده می‌کند.

تفاوت Direct و Topic Exchange چیست؟

Direct Exchange تطبیق دقیق Routing Key انجام می‌دهد. Topic Exchange Routing Key را با Patternهای شامل * و # مقایسه می‌کند.

Ack چیست؟

Ack تأییدی است که Consumer پس از پردازش موفق Message برای RabbitMQ ارسال می‌کند. بعد از Ack، Message از Queue حذف می‌شود.

Prefetch چیست؟

Prefetch حداکثر تعداد Messageهای تأییدنشده‌ای است که RabbitMQ هم‌زمان به Consumer تحویل می‌دهد.

DLQ چیست؟

Dead Letter Queue محلی برای Messageهایی است که قابل پردازش نبوده‌اند، منقضی شده‌اند یا با requeue=false رد شده‌اند.

RabbitMQ بهتر است یا Redis؟

برای Queue ساده و پروژه‌ای که از قبل Redis دارد، Redis می‌تواند کافی باشد. برای Routing پیشرفته، Ack، DLQ و مدیریت تخصصی Message، RabbitMQ مناسب‌تر است.

RabbitMQ بهتر است یا Kafka؟

RabbitMQ برای Task Queue و Routing پیام مناسب است. Kafka برای Event Streaming، Retention و Replay جریان‌های حجیم طراحی شده است. انتخاب به نیاز پروژه بستگی دارد.

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

بله. کتابخانه aio-pika امکان اتصال Async به RabbitMQ را فراهم می‌کند. بهتر است Connection در Lifespan برنامه ساخته شود، نه در هر Request.

آیا RabbitMQ از پردازش تکراری جلوگیری می‌کند؟

به‌طور مطلق خیر. در بعضی شرایط Message دوباره تحویل داده می‌شود. Consumer باید با event_id یا job_id به‌صورت Idempotent طراحی شود.

آیا RabbitMQ برای پردازش هوش مصنوعی مناسب است؟

بله. پردازش‌های هوش مصنوعی معمولاً کندتر از APIهای معمولی هستند و می‌توانند در Workerهای جداگانه اجرا شوند. RabbitMQ درخواست‌ها را میان Workerها توزیع می‌کند.

هزینه پردازش هوش مصنوعی چگونه محاسبه می‌شود؟

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

جمع‌بندی

RabbitMQ یک Message Broker قدرتمند برای جداکردن عملیات سنگین از مسیر اصلی API و ساخت سیستم‌های Message-Driven است. مفاهیمی مانند Exchange، Queue، Binding، Routing Key، Ack، Publisher Confirm، Prefetch و Dead Letter Exchange پایه‌های اصلی کار با RabbitMQ هستند.

در پروژه عملی این مقاله، RabbitMQ را با Docker راه‌اندازی کردیم و یک Pipeline کامل با FastAPI، Python، aio-pika و API درواره ساختیم. FastAPI درخواست کاربر را به Job تبدیل می‌کند، RabbitMQ آن را در Queue قرار می‌دهد و Worker مستقل متن را با مدل هوش مصنوعی پردازش می‌کند. نتیجه در Queue جداگانه منتشر می‌شود و Jobهای ناموفق نیز پس از Retry محدود به DLQ منتقل می‌شوند.

این معماری برای خلاصه‌سازی، ترجمه، دسته‌بندی، تحلیل سند، تولید محتوا، پردازش تصویر و سایر Workflowهای غیرهم‌زمان قابل توسعه است. بااین‌حال، در سیستم‌های مهم باید Idempotency، Outbox Pattern، Monitoring، Publisher Confirm، DLQ و بازیابی پس از خطا از ابتدا در طراحی در نظر گرفته شوند.

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

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

مقالات مرتبط

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

Read more

اتوماسیون هوش مصنوعی چیست؟ کاربردها و آموزش ساخت AI Automation

اتوماسیون هوش مصنوعی چیست؟ کاربردها و آموزش ساخت AI Automation

اتوماسیون هوش مصنوعی با ترکیب گردش‌کارهای خودکار و مدل‌های هوش مصنوعی، پردازش متن، دسته‌بندی، استخراج اطلاعات و تصمیم‌های پیشنهادی را خودکار می‌کند. در این راهنما، معماری و ساخت نمونه عملی آن با API درواره را می‌آموزید.

Agentic Commerce چیست؟ آینده خرید با ایجنت هوش مصنوعی

Agentic Commerce چیست؟ آینده خرید با ایجنت هوش مصنوعی

Agentic Commerce شیوه‌ای جدید برای خرید اینترنتی است که در آن ایجنت هوش مصنوعی می‌تواند نیاز کاربر را بفهمد، محصولات را جست‌وجو و مقایسه کند و فرایند خرید را پیش ببرد. در این راهنما با معماری، UCP، ACP و پیاده‌سازی آن با API درواره آشنا می‌شوید.