RabbitMQ چیست؟ آموزش کامل Message Queue با Python، FastAPI و API درواره
RabbitMQ یک Message Broker برای اجرای Jobهای غیرهمزمان و ارتباط میان سرویسهاست. در این آموزش، RabbitMQ را با Docker اجرا میکنید و یک سیستم عملی با Python، FastAPI، Retry، DLQ و 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_idcorrelation_idcontent_typecontent_encodingtimestampheadersdelivery_modeexpirationreply_totype
نمونه 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 تمرکز دارد.
| ویژگی | RabbitMQ | Redis |
|---|---|---|
| کاربرد اصلی | Message Broker | In-Memory Data Store |
| Routing | Direct، Topic، Fanout و Headers | محدودتر |
| Acknowledgement | قابلیت اصلی | وابسته به ساختار انتخابی |
| Dead Letter | پشتیبانی دارد | نیازمند طراحی |
| Queue Management | تخصصی | یکی از کاربردها |
| Cache | کاربرد اصلی نیست | بسیار مناسب |
| Message Priority | پشتیبانی در Classic Queue | نیازمند طراحی |
| Pub/Sub | دارد | دارد |
| Persistent Messaging | دارد | وابسته به Persistence |
| Management UI | Plugin رسمی دارد | ابزارهای جداگانه |
اگر پروژه از قبل Redis دارد و Queue سادهای نیاز دارد، Redis ممکن است کافی باشد. اگر Routing، Ack، DLQ و رفتار دقیق Message Broker اهمیت دارد، RabbitMQ انتخاب تخصصیتری است.
RabbitMQ چه تفاوتی با Kafka دارد؟
RabbitMQ و Kafka در بعضی سناریوها همپوشانی دارند، اما مدل اصلی آنها متفاوت است.
| معیار | RabbitMQ | Kafka |
|---|---|---|
| مدل اصلی | Message Broker و Queue | Distributed 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:
- Message با Manual Ack دریافت میشود.
- متن حداکثر سه بار پردازش میشود.
- نتیجه در Exchange خروجی منتشر میشود.
- Publisher Confirm نتیجه بررسی میشود.
- فقط پس از خروج موفق از Context Manager، Message اصلی Ack میشود.
- اگر پردازش شکست بخورد، 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 را دوباره تحویل دهد. برای مثال:
- Worker Message را پردازش میکند.
- نتیجه را در Database ذخیره میکند.
- پیش از Ack متوقف میشود.
- RabbitMQ Message را به Worker دیگری میدهد.
- عملیات دوباره اجرا میشود.
بنابراین هر 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 این است:
- داده در Database ذخیره شود.
- 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 بسازید و مدل مناسب پروژه را از صفحه مدلها انتخاب کنید.
منابع پیشنهادی
- مستندات رسمی RabbitMQ
- راهنمای نصب RabbitMQ
- آموزش Work Queue با Python
- مستندات Exchangeها
- مستندات Consumer Ack و Publisher Confirm
- مستندات Dead Letter Exchange
- مستندات aio-pika
- مستندات FastAPI
مقالات مرتبط
- پردازش غیرهمزمان API هوش مصنوعی با Celery، Redis، Worker و Webhook
- ساخت API هوش مصنوعی آماده Production
- مانیتورینگ و Observability در سیستمهای هوش مصنوعی
- آموزش اتصال API هوش مصنوعی به اپلیکیشن
- تست Load و Performance API با k6 و هوش مصنوعی
- دستهبندی تیکتهای پشتیبانی با هوش مصنوعی
برای مطالعه شرایط استفاده و محدودیتهای مسئولیت، صفحه «سلب مسئولیت» را مشاهده کنید.