Kafka چیست؟ آموزش کامل Apache Kafka، Event Streaming و اتصال به Python و FastAPI
Apache Kafka پلتفرمی برای انتقال و پردازش جریان رویدادهاست. در این آموزش، Kafka را با Docker راهاندازی میکنید و یک سیستم Event-Driven عملی با Python، FastAPI و API هوش مصنوعی درواره میسازید.
وقتی یک برنامه کوچک است، معمولاً همه عملیات در همان درخواست HTTP انجام میشوند. کاربر درخواستی ارسال میکند، Backend داده را پردازش میکند، چند API را فراخوانی میکند و پس از پایان همه مراحل پاسخ را برمیگرداند.
این معماری در شروع ساده و قابلفهم است، اما با افزایش تعداد کاربران و سرویسها مشکلاتی ایجاد میکند:
- پاسخ API به عملیات کند وابسته میشود.
- اختلال یک سرویس میتواند کل زنجیره را متوقف کند.
- پردازش همزمان هزاران درخواست دشوار میشود.
- اضافهکردن مصرفکننده جدید به تغییر سرویس اصلی نیاز دارد.
- رویدادهای گذشته بهراحتی قابل پردازش مجدد نیستند.
- کنترل Retry، ترتیب پیامها و توزیع بار پیچیده میشود.
Apache Kafka یا آپاچی کافکا یکی از مهمترین ابزارهای معماری Event-Driven و Event Streaming است. Kafka رویدادها را دریافت، ذخیره و میان Producerها و Consumerها توزیع میکند. چند سرویس مستقل میتوانند یک جریان رویداد را با سرعت و مقیاس بالا پردازش کنند، بدون آنکه Producer مستقیماً به همه آنها وابسته باشد.
در این مقاله Kafka را از پایه یاد میگیریم، آن را با Docker اجرا میکنیم و سپس یک سیستم واقعی میسازیم که درخواست پردازش متن را از FastAPI دریافت میکند، آن را در Kafka قرار میدهد و با یک Worker جداگانه از طریق API درواره پردازش میکند.
Kafka چیست؟
Apache Kafka یک پلتفرم توزیعشده Event Streaming است. برنامهها میتوانند رویدادها را در Kafka بنویسند، بخوانند، ذخیره و پردازش کنند.
هر رویداد میتواند نشاندهنده اتفاقی در سیستم باشد:
- کاربر ثبتنام کرد.
- سفارش ایجاد شد.
- پرداخت تأیید شد.
- فایل بارگذاری شد.
- متن برای خلاصهسازی ارسال شد.
- مدل هوش مصنوعی پاسخ داد.
- موجودی محصول تغییر کرد.
- گزارش آماده شد.
- Sensor مقدار جدیدی ارسال کرد.
- وضعیت یک Job تغییر کرد.
براساس مستندات رسمی Apache Kafka، برنامههای Producer رویدادها را منتشر میکنند و Consumerها آنها را میخوانند و پردازش میکنند. Kafka این رویدادها را در Topicها نگه میدارد و میتواند پردازش آنها را میان چند Consumer توزیع کند.
Event Streaming چیست؟
Event Streaming یعنی ثبت، انتقال و پردازش پیوسته رویدادهایی که در یک سیستم رخ میدهند.
در معماری Request/Response سنتی، سرویس A مستقیماً سرویس B را فراخوانی میکند:
API → سرویس پردازش → پاسخ
در معماری Event-Driven، سرویس A یک Event منتشر میکند و لزوماً منتظر پایان پردازش نمیماند:
API → Kafka Topic → Worker
سرویسهای دیگری نیز میتوانند همان جریان را مصرف کنند:
Kafka Topic
├── AI Worker
├── Analytics Service
├── Notification Service
└── Audit Service
مزیت مهم این معماری Decoupling یا کاهش وابستگی مستقیم سرویسهاست. Producer فقط Event را منتشر میکند و لازم نیست همه Consumerها را بشناسد.
Kafka چه مشکلی را حل میکند؟
فرض کنید کاربر متنی را برای تحلیل هوش مصنوعی ارسال میکند. انجام همه مراحل در یک درخواست ممکن است شامل این عملیات باشد:
- اعتبارسنجی متن
- ذخیره درخواست
- ارسال به مدل هوش مصنوعی
- دریافت خروجی
- ذخیره نتیجه
- ارسال Notification
- ثبت آمار مصرف
- بهروزرسانی داشبورد
اگر همه این مراحل در یک HTTP Request انجام شوند، زمان پاسخ افزایش پیدا میکند. اختلال موقت یکی از سرویسها نیز ممکن است باعث شکست کل درخواست شود.
با Kafka میتوان درخواست را به Event تبدیل کرد:
{
"event_type": "ai.text.requested",
"job_id": "job-123",
"payload": {
"text": "متن موردنظر کاربر",
"operation": "summarize"
}
}
FastAPI رویداد را در Kafka مینویسد و سریع شناسه Job را برمیگرداند. Worker مستقل بعداً Event را دریافت و پردازش میکند.
مفاهیم اصلی Apache Kafka
برای کار با Kafka باید چند مفهوم پایه را دقیق بشناسیم.
Event یا Message
Event یک رکورد داده است که وقوع یک اتفاق را نشان میدهد. هر Event در Kafka معمولاً شامل این بخشهاست:
- Key
- Value
- Timestamp
- Header
- Topic
- Partition
- Offset
نمونه Value:
{
"event_id": "0ae5d826-7f1d-44ae-97e4-a53055cab051",
"event_type": "ai.text.requested",
"event_version": 1,
"occurred_at": "2026-08-06T12:00:00Z",
"payload": {
"text": "این متن را خلاصه کن",
"operation": "summarize"
}
}
Producer
Producer برنامهای است که Event را در Kafka منتشر میکند.
در پروژه این مقاله، FastAPI نقش Producer را دارد:
FastAPI → Kafka
Consumer
Consumer برنامهای است که Eventها را از Kafka میخواند و پردازش میکند.
در پروژه ما، AI Worker یک Consumer است:
Kafka → AI Worker → API درواره
Broker
هر نمونه در حال اجرای Kafka یک Broker نامیده میشود. در محیط توسعه میتوان یک Broker داشت، اما Clusterهای Production معمولاً از چند Broker تشکیل میشوند.
Cluster
مجموعه Brokerهایی که با هم کار میکنند یک Kafka Cluster را تشکیل میدهند.
Topic
Topic جریان مشخصی از Eventها را نگه میدارد. هر نوع Event بهتر است در Topic متناسب قرار بگیرد.
نمونه Topicها:
ai.text.requested.v1
ai.text.completed.v1
ai.text.failed.v1
order.created.v1
payment.completed.v1
notification.requested.v1
نامگذاری واضح Topicها نگهداری سامانه را سادهتر میکند.
Partition
هر Topic میتواند به چند Partition تقسیم شود. Partitionها امکان توزیع Storage و پردازش موازی را فراهم میکنند.
فرض کنید Topic زیر سه Partition دارد:
ai.text.requested.v1
├── Partition 0
├── Partition 1
└── Partition 2
Eventها میان این Partitionها توزیع میشوند. هر Partition یک Log مرتب و افزایشی است.
Offset
هر Event در یک Partition شمارهای به نام Offset دارد:
Partition 0:
Offset 0 → Event A
Offset 1 → Event B
Offset 2 → Event C
Consumer با Offset متوجه میشود تا کجای جریان را پردازش کرده است.
Offset در کل Topic یکتا نیست؛ بلکه در هر Partition معنا دارد. بنابراین یک موقعیت کامل شامل Topic، Partition و Offset است.
Consumer Group
Consumer Group مجموعهای از Consumerهاست که برای پردازش یک یا چند Topic با یکدیگر همکاری میکنند.
اگر Topic سه Partition داشته باشد و سه Consumer در یک Group باشند، هر Consumer میتواند یک Partition را پردازش کند:
Partition 0 → Consumer 1
Partition 1 → Consumer 2
Partition 2 → Consumer 3
در هر Consumer Group، یک Partition در یک زمان معمولاً فقط به یک Consumer اختصاص داده میشود. بنابراین اگر Topic سه Partition داشته باشد و پنج Consumer در همان Group اجرا شوند، دو Consumer ممکن است بیکار بمانند.
اما اگر دو Consumer Group متفاوت داشته باشیم، هر Group جریان را مستقل دریافت میکند:
Topic
├── Consumer Group: ai-workers
└── Consumer Group: analytics
در این حالت هم Workerهای هوش مصنوعی و هم سرویس Analytics میتوانند همه Eventها را مستقل پردازش کنند.
Retention
Kafka برخلاف یک Queue ساده لزوماً Event را بلافاصله پس از مصرف حذف نمیکند. Eventها براساس سیاست Retention برای مدت مشخص یا تا رسیدن Log به حجم تعیینشده نگهداری میشوند.
این ویژگی اجازه میدهد:
- Consumer جدید Eventهای گذشته را بخواند.
- دادهها دوباره پردازش شوند.
- Pipeline جدید روی Eventهای تاریخی اجرا شود.
- Consumer پس از قطعی از Offset قبلی ادامه دهد.
Replica
برای افزایش دسترسپذیری، هر Partition میتواند روی چند Broker Replica داشته باشد. یکی از Replicaها Leader است و سایر Replicaها داده را دنبال میکنند.
در محیط تکBroker توسعه، Replication Factor برابر ۱ خواهد بود. در Production باید براساس تعداد Broker و سطح دسترسپذیری موردنیاز تنظیم شود.
Kafka از ZooKeeper استفاده میکند؟
نسخههای جدید Kafka از معماری KRaft استفاده میکنند و برای مدیریت Metadata به ZooKeeper نیاز ندارند. در این آموزش از Apache Kafka 4.3.1 استفاده میکنیم که براساس Quickstart رسمی در حالت KRaft اجرا میشود.
نسخه جاری و دستور اجرای رسمی در صفحه Kafka Quickstart قابل بررسی است.
Kafka چه تفاوتی با Message Queue دارد؟
Kafka میتواند برای بعضی سناریوهای صف استفاده شود، اما مدل آن با Queueهای سنتی یکسان نیست.
| ویژگی | Apache Kafka | Queue سنتی |
|---|---|---|
| مدل اصلی | Event Log توزیعشده | صف پیام |
| نگهداری پس از مصرف | براساس Retention | معمولاً حذف پس از Ack |
| Replay | قابلیت اصلی | معمولاً محدودتر |
| ترتیب | داخل هر Partition | وابسته به Queue |
| پردازش موازی | با Partition و Consumer Group | با Workerهای صف |
| چند Consumer مستقل | با Groupهای مختلف | بسته به ابزار و Exchange |
| کاربرد مناسب | Event Streaming، Integration، Analytics | Task Queue و Job Processing |
| مقیاس جریان داده | بسیار بالا | وابسته به محصول |
انتخاب Kafka فقط بهدلیل محبوبیت آن تصمیم مناسبی نیست. اگر صرفاً چند Job پسزمینه ساده دارید، ابزارهایی مانند Redis Queue، RabbitMQ یا Celery ممکن است سادهتر باشند.
مقایسه Kafka، RabbitMQ و Redis Streams
| معیار | Kafka | RabbitMQ | Redis Streams |
|---|---|---|---|
| کاربرد اصلی | Event Streaming | Message Broker | Stream سبک در Redis |
| Replay رویداد | بسیار مناسب | محدودتر | امکانپذیر |
| Throughput | بالا | مناسب | بالا در سناریوهای سبک |
| Routing پیچیده | متوسط | بسیار خوب | محدودتر |
| Retention طولانی | مناسب | کاربرد اصلی نیست | وابسته به حافظه و تنظیمات |
| پیچیدگی عملیاتی | بیشتر | متوسط | کمتر در صورت وجود Redis |
| Consumer Group | دارد | Worker Consumer دارد | دارد |
| مناسب Event Sourcing | مناسبتر | کمتر | برای پروژههای کوچک |
| مناسب Task Queue ساده | ممکن است بیشازحد پیچیده باشد | مناسب | مناسب |
چه زمانی از Kafka استفاده کنیم؟
Kafka معمولاً در این سناریوها ارزش ایجاد میکند:
- چند سرویس باید یک Event را مستقل مصرف کنند.
- Eventها باید برای پردازش مجدد نگهداری شوند.
- حجم بالایی از رویدادها دارید.
- ترتیب رویدادهای مرتبط اهمیت دارد.
- معماری Microservices دارید.
- پردازش Real-Time یا Near Real-Time نیاز دارید.
- چند Pipeline تحلیلی از یک جریان مشترک استفاده میکنند.
- میخواهید Producer و Consumer مستقل مقیاسپذیر باشند.
- دادهها از چند منبع وارد سامانه میشوند.
- Workflowهای هوش مصنوعی غیرهمزمان دارید.
چه زمانی Kafka انتخاب مناسبی نیست؟
Kafka احتمالاً برای این سناریوها انتخاب اول نیست:
- پروژه کوچک و کمترافیک است.
- تنها یک Worker ساده دارید.
- Replay رویدادها لازم نیست.
- تیم تجربه مدیریت Cluster ندارد.
- عملیات باید کاملاً Synchronous باشد.
- پیچیدگی جدید ارزش تجاری مشخصی ایجاد نمیکند.
- یک Queue سبک تمام نیاز را پوشش میدهد.
برای MVP بهتر است سادهترین معماری پاسخگو به نیاز واقعی انتخاب شود.
تضمین ترتیب Eventها در Kafka
Kafka ترتیب را فقط داخل یک Partition تضمین میکند، نه میان تمام Partitionهای یک Topic.
اگر ترتیب Eventهای یک کاربر مهم است، باید Eventهای همان کاربر با Key یکسان منتشر شوند:
Key = user-123
Kafka معمولاً Eventهای دارای Key یکسان را به یک Partition میفرستد. در نتیجه ترتیب Eventهای آن Key داخل همان Partition حفظ میشود.
نمونه:
user-123:
1. profile.created
2. profile.updated
3. profile.deleted
اگر Eventها بدون Key یا با Keyهای متفاوت ارسال شوند، ممکن است در Partitionهای مختلف قرار بگیرند و ترتیب سراسری قابل تضمین نباشد.
Delivery Semantics در Kafka
سه اصطلاح مهم در پردازش پیام وجود دارد.
At-most-once
هر Event حداکثر یکبار پردازش میشود، اما ممکن است در صورت خطا از دست برود.
At-least-once
هر Event حداقل یکبار پردازش میشود، اما امکان پردازش تکراری وجود دارد.
این مدل برای بسیاری از Consumerها رایج است. برنامه باید Idempotent باشد تا پردازش دوباره Event نتیجه نامعتبر ایجاد نکند.
Exactly-once
Kafka قابلیتهای Transactional برای پردازش Exactly-once در محدوده مشخص Kafka ارائه میکند، اما این اصطلاح نباید بیشازحد تعمیم داده شود.
اگر Consumer یک API خارجی، درگاه پرداخت، ایمیل یا مدل هوش مصنوعی را فراخوانی کند، Exactly-once سراسری بهطور خودکار ایجاد نمیشود. در این شرایط همچنان به Idempotency Key، ذخیره وضعیت و کنترل Side Effect نیاز دارید.
Idempotency چیست؟
پردازش Idempotent یعنی اجرای چندباره یک Event، اثر نهایی ناخواسته ایجاد نکند.
هر Event باید یک شناسه یکتا داشته باشد:
{
"event_id": "93ea4764-f887-4a3c-829e-588b4ca5bfaa"
}
Consumer پیش از پردازش میتواند بررسی کند آیا این event_id قبلاً تکمیل شده است یا خیر.
الگوی ساده:
دریافت Event
→ بررسی event_id
→ اگر قبلاً پردازش شده: ردکردن
→ اگر جدید است: پردازش
→ ثبت event_id
→ Commit Offset
در Production بهتر است ثبت نتیجه پردازش و وضعیت Idempotency در یک Storage پایدار انجام شود.
نصب Kafka با Docker
برای اجرای محلی Kafka به Docker و Docker Compose نیاز داریم.
ساختار پروژه:
kafka-ai-pipeline/
├── docker-compose.yml
├── requirements.txt
├── .env
├── producer_api.py
├── worker.py
└── result_consumer.py
فایل docker-compose.yml را بسازید:
services:
kafka:
image: apache/kafka:4.3.1
container_name: darvareh-kafka
ports:
- "127.0.0.1:9092:9092"
volumes:
- kafka_data:/var/lib/kafka/data
volumes:
kafka_data:
Kafka را اجرا کنید:
docker compose up -d
وضعیت Container:
docker compose ps
مشاهده Log:
docker compose logs -f kafka
این پیکربندی برای آموزش و محیط محلی است. برای Production باید تنظیمات Cluster چندBroker، Replication، Storage، دسترسی شبکه، احراز هویت، TLS، Monitoring و Backup متناسب با زیرساخت شما طراحی شود.
ساخت اولین Topic
برای ساخت Topic درخواستهای پردازش متن:
docker exec darvareh-kafka \
/opt/kafka/bin/kafka-topics.sh \
--create \
--topic ai.text.requested.v1 \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
Topic نتایج موفق:
docker exec darvareh-kafka \
/opt/kafka/bin/kafka-topics.sh \
--create \
--topic ai.text.completed.v1 \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
Topic خطا یا Dead Letter Topic:
docker exec darvareh-kafka \
/opt/kafka/bin/kafka-topics.sh \
--create \
--topic ai.text.failed.v1 \
--bootstrap-server localhost:9092 \
--partitions 3 \
--replication-factor 1
مشاهده فهرست Topicها:
docker exec darvareh-kafka \
/opt/kafka/bin/kafka-topics.sh \
--list \
--bootstrap-server localhost:9092
مشاهده مشخصات Topic:
docker exec darvareh-kafka \
/opt/kafka/bin/kafka-topics.sh \
--describe \
--topic ai.text.requested.v1 \
--bootstrap-server localhost:9092
تست Kafka با Console Producer
Producer خط فرمان را اجرا کنید:
docker exec -it darvareh-kafka \
/opt/kafka/bin/kafka-console-producer.sh \
--topic ai.text.requested.v1 \
--bootstrap-server localhost:9092
سپس یک پیام JSON وارد کنید:
{"event_type":"ai.text.requested","payload":{"text":"این متن را خلاصه کن"}}
خواندن Event با Console Consumer
در Terminal دیگری اجرا کنید:
docker exec -it darvareh-kafka \
/opt/kafka/bin/kafka-console-consumer.sh \
--topic ai.text.requested.v1 \
--bootstrap-server localhost:9092 \
--from-beginning
Event ارسالشده باید نمایش داده شود.
ساخت پروژه عملی Kafka با Python و FastAPI
در این پروژه سه برنامه داریم:
producer_api.pyدرخواست HTTP را دریافت و Event را منتشر میکند.worker.pyEvent را مصرف و متن را از طریق API درواره پردازش میکند.result_consumer.pyنتیجه پردازش را از Topic خروجی میخواند.
نصب کتابخانهها
فایل requirements.txt:
fastapi>=0.115,<1
uvicorn[standard]>=0.34,<1
confluent-kafka>=2.15,<3
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
کتابخانه confluent-kafka یک Client سریع مبتنی بر librdkafka برای Python ارائه میکند. جزئیات نصب و API آن در مستندات Python Client for Apache Kafka در دسترس است.
تنظیم متغیرهای محیطی
فایل .env:
KAFKA_BOOTSTRAP_SERVERS=localhost:9092
KAFKA_REQUEST_TOPIC=ai.text.requested.v1
KAFKA_COMPLETED_TOPIC=ai.text.completed.v1
KAFKA_FAILED_TOPIC=ai.text.failed.v1
DARVAREH_API_KEY=YOUR_DARVAREH_API_KEY
DARVAREH_MODEL_ID=YOUR_MODEL_ID
DARVAREH_BASE_URL=https://api.darvareh.ir/v1
فایل .env را در Git ثبت نکنید:
.env
.venv/
__pycache__/
ساخت Producer با FastAPI
فایل producer_api.py:
import json
import os
from contextlib import asynccontextmanager
from datetime import datetime, timezone
from typing import Literal
from uuid import uuid4
from confluent_kafka import KafkaException, Producer
from dotenv import load_dotenv
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel, Field
load_dotenv()
BOOTSTRAP_SERVERS = os.getenv(
"KAFKA_BOOTSTRAP_SERVERS",
"localhost:9092",
)
REQUEST_TOPIC = os.getenv(
"KAFKA_REQUEST_TOPIC",
"ai.text.requested.v1",
)
producer = Producer(
{
"bootstrap.servers": BOOTSTRAP_SERVERS,
"client.id": "ai-job-api",
"acks": "all",
"enable.idempotence": True,
"compression.type": "snappy",
"linger.ms": 10,
}
)
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()
def delivery_report(error, message) -> None:
if error is not None:
print(
"Event delivery failed:",
error,
)
return
print(
"Event delivered:",
{
"topic": message.topic(),
"partition": message.partition(),
"offset": message.offset(),
},
)
@asynccontextmanager
async def lifespan(app: FastAPI):
yield
producer.flush(10)
app = FastAPI(
title="Kafka AI Job API",
version="1.0.0",
lifespan=lifespan,
)
@app.get("/health")
def health():
return {
"status": "ok",
"kafka": BOOTSTRAP_SERVERS,
}
@app.post("/jobs", status_code=202)
def create_job(request: AIJobRequest):
job_id = str(uuid4())
event_id = str(uuid4())
event = {
"event_id": event_id,
"event_type": "ai.text.requested",
"event_version": 1,
"occurred_at": utc_now(),
"trace_id": job_id,
"payload": {
"job_id": job_id,
"user_id": request.user_id,
"operation": request.operation,
"text": request.text,
},
}
try:
producer.produce(
topic=REQUEST_TOPIC,
key=job_id.encode("utf-8"),
value=json.dumps(
event,
ensure_ascii=False,
).encode("utf-8"),
headers={
"event_type": "ai.text.requested",
"event_version": "1",
"trace_id": job_id,
},
callback=delivery_report,
)
producer.poll(0)
except BufferError as exc:
raise HTTPException(
status_code=503,
detail="Kafka producer queue is full",
) from exc
except KafkaException as exc:
raise HTTPException(
status_code=503,
detail="Event could not be queued",
) from exc
return {
"job_id": job_id,
"event_id": event_id,
"status": "queued",
}
چند نکته مهم در Producer:
acks=allباعث میشود Producer تأیید همه Replicaهای همگام موردنیاز را دریافت کند.enable.idempotence=trueاحتمال ثبت تکراری ناشی از Retry داخلی Producer را کاهش میدهد.key=job_idباعث میشود Eventهای یک Job به یک Partition بروند.- پاسخ HTTP برابر
202 Acceptedاست، زیرا پردازش هنوز تکمیل نشده است. producer.poll(0)Callbackهای Delivery را پردازش میکند.- هنگام خاموششدن برنامه،
flushبرای ارسال Eventهای باقیمانده فراخوانی میشود.
در معماری Production، پذیرش قطعی Job بهتر است با Outbox Pattern یا سازوکار پایدار مشابه طراحی شود. صرفاً برگرداندن 202 پیش از تأیید پایدار Event ممکن است در شرایط خاص باعث ابهام وضعیت درخواست شود.
ساخت AI Worker
فایل worker.py:
import json
import os
import signal
import time
from datetime import datetime, timezone
from uuid import uuid4
from confluent_kafka import Consumer, KafkaException, Producer
from dotenv import load_dotenv
from openai import OpenAI
load_dotenv()
BOOTSTRAP_SERVERS = os.getenv(
"KAFKA_BOOTSTRAP_SERVERS",
"localhost:9092",
)
REQUEST_TOPIC = os.getenv(
"KAFKA_REQUEST_TOPIC",
"ai.text.requested.v1",
)
COMPLETED_TOPIC = os.getenv(
"KAFKA_COMPLETED_TOPIC",
"ai.text.completed.v1",
)
FAILED_TOPIC = os.getenv(
"KAFKA_FAILED_TOPIC",
"ai.text.failed.v1",
)
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"
)
consumer = Consumer(
{
"bootstrap.servers": BOOTSTRAP_SERVERS,
"group.id": "ai-text-workers-v1",
"client.id": "ai-text-worker",
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
"max.poll.interval.ms": 300_000,
}
)
producer = Producer(
{
"bootstrap.servers": BOOTSTRAP_SERVERS,
"client.id": "ai-result-producer",
"acks": "all",
"enable.idempotence": True,
"compression.type": "snappy",
}
)
ai_client = OpenAI(
api_key=DARVAREH_API_KEY,
base_url=DARVAREH_BASE_URL,
timeout=60,
max_retries=2,
)
running = True
def stop_worker(signum, frame):
global running
running = False
signal.signal(signal.SIGINT, stop_worker)
signal.signal(signal.SIGTERM, stop_worker)
def utc_now() -> str:
return datetime.now(timezone.utc).isoformat()
def build_messages(
operation: str,
text: str,
) -> list[dict[str, str]]:
prompts = {
"summarize": (
"متن کاربر را دقیق و کوتاه خلاصه کن. "
"هیچ اطلاعاتی به متن اضافه نکن."
),
"rewrite": (
"متن کاربر را به فارسی روان و حرفهای "
"بازنویسی کن، بدون تغییر مفهوم."
),
"classify": (
"موضوع اصلی متن را مشخص کن. "
"خروجی را فقط در یک عبارت کوتاه بنویس."
),
}
instruction = prompts.get(
operation,
prompts["summarize"],
)
return [
{
"role": "system",
"content": instruction,
},
{
"role": "user",
"content": text,
},
]
def call_ai_with_retry(
operation: str,
text: str,
max_attempts: int = 3,
) -> str:
last_error = None
for attempt in range(1, max_attempts + 1):
try:
response = ai_client.chat.completions.create(
model=DARVAREH_MODEL_ID,
temperature=0.2,
max_tokens=800,
messages=build_messages(
operation=operation,
text=text,
),
)
content = response.choices[
0
].message.content
if not content:
raise ValueError(
"AI response content is empty"
)
return content
except Exception as exc:
last_error = exc
if attempt < max_attempts:
delay = 2 ** (attempt - 1)
time.sleep(delay)
raise RuntimeError(
"AI processing failed after retries"
) from last_error
def publish_event(
topic: str,
key: str,
event: dict,
) -> None:
producer.produce(
topic=topic,
key=key.encode("utf-8"),
value=json.dumps(
event,
ensure_ascii=False,
).encode("utf-8"),
headers={
"event_type": event["event_type"],
"event_version": str(
event["event_version"]
),
"trace_id": event["trace_id"],
},
)
remaining = producer.flush(10)
if remaining > 0:
raise RuntimeError(
"Result event was not delivered"
)
def process_message(message) -> None:
source_event = json.loads(
message.value().decode("utf-8")
)
payload = source_event["payload"]
job_id = payload["job_id"]
operation = payload["operation"]
text = payload["text"]
try:
result = call_ai_with_retry(
operation=operation,
text=text,
)
completed_event = {
"event_id": str(uuid4()),
"event_type": "ai.text.completed",
"event_version": 1,
"occurred_at": utc_now(),
"trace_id": source_event["trace_id"],
"causation_id": source_event["event_id"],
"payload": {
"job_id": job_id,
"operation": operation,
"result": result,
"model": DARVAREH_MODEL_ID,
},
}
publish_event(
topic=COMPLETED_TOPIC,
key=job_id,
event=completed_event,
)
consumer.commit(
message=message,
asynchronous=False,
)
print(
f"Job {job_id} completed"
)
except Exception as exc:
failed_event = {
"event_id": str(uuid4()),
"event_type": "ai.text.failed",
"event_version": 1,
"occurred_at": utc_now(),
"trace_id": source_event["trace_id"],
"causation_id": source_event["event_id"],
"payload": {
"job_id": job_id,
"operation": operation,
"error_type": type(exc).__name__,
"error_message": str(exc)[:500],
},
}
publish_event(
topic=FAILED_TOPIC,
key=job_id,
event=failed_event,
)
consumer.commit(
message=message,
asynchronous=False,
)
print(
f"Job {job_id} moved to failed topic"
)
def main() -> None:
consumer.subscribe([REQUEST_TOPIC])
print(
f"Worker is consuming {REQUEST_TOPIC}"
)
try:
while running:
message = consumer.poll(1.0)
if message is None:
continue
if message.error():
raise KafkaException(
message.error()
)
process_message(message)
finally:
consumer.close()
producer.flush(10)
ai_client.close()
print("Worker stopped")
if __name__ == "__main__":
main()
در این Worker، Offset فقط بعد از انتشار موفق Event نتیجه یا Event خطا Commit میشود. این ترتیب احتمال گمشدن خاموش Event را کاهش میدهد:
دریافت درخواست
→ پردازش با مدل
→ انتشار نتیجه
→ تأیید انتشار
→ Commit Offset
بااینحال، اگر Worker پس از انتشار نتیجه و پیش از Commit متوقف شود، ممکن است Event دوباره پردازش شود. برای کنترل کامل باید Consumer خروجی و Storage نتیجه براساس event_id یا job_id رفتار Idempotent داشته باشند.
ساخت Consumer نتایج
فایل result_consumer.py:
import json
import os
import signal
from confluent_kafka import Consumer, KafkaException
from dotenv import load_dotenv
load_dotenv()
BOOTSTRAP_SERVERS = os.getenv(
"KAFKA_BOOTSTRAP_SERVERS",
"localhost:9092",
)
COMPLETED_TOPIC = os.getenv(
"KAFKA_COMPLETED_TOPIC",
"ai.text.completed.v1",
)
consumer = Consumer(
{
"bootstrap.servers": BOOTSTRAP_SERVERS,
"group.id": "ai-result-readers-v1",
"client.id": "ai-result-reader",
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
}
)
running = True
def stop_consumer(signum, frame):
global running
running = False
signal.signal(signal.SIGINT, stop_consumer)
signal.signal(signal.SIGTERM, stop_consumer)
def main() -> None:
consumer.subscribe([COMPLETED_TOPIC])
print(
f"Reading results from {COMPLETED_TOPIC}"
)
try:
while running:
message = consumer.poll(1.0)
if message is None:
continue
if message.error():
raise KafkaException(
message.error()
)
event = json.loads(
message.value().decode("utf-8")
)
payload = event["payload"]
print(
"\nJob:",
payload["job_id"],
)
print(
"Operation:",
payload["operation"],
)
print(
"Result:",
payload["result"],
)
consumer.commit(
message=message,
asynchronous=False,
)
finally:
consumer.close()
print("Result consumer stopped")
if __name__ == "__main__":
main()
اجرای کامل پروژه
ابتدا Kafka را اجرا کنید:
docker compose up -d
Topicها را مطابق دستورهای قبلی بسازید.
سپس FastAPI را اجرا کنید:
uvicorn producer_api:app --reload
در Terminal دوم Worker را اجرا کنید:
python worker.py
در Terminal سوم Consumer نتایج را اجرا کنید:
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"
}'
پاسخ FastAPI:
{
"job_id": "d307fe67-944a-4914-8699-7a220c398da3",
"event_id": "95173f3e-d5d1-486d-b2e0-2244623e7aac",
"status": "queued"
}
پس از پردازش، نتیجه در Terminal مربوط به result_consumer.py نمایش داده میشود.
برای ساخت کلید API و استفاده از مدلهای مختلف میتوانید وارد درواره شوید. فهرست مدلهای فعال، شناسه مدلها و قیمت بهروز آنها در صفحه مدلهای درواره قرار دارد.
چرا FastAPI مستقیماً منتظر مدل نمیماند؟
در معماری این پروژه، FastAPI فقط Job را میپذیرد و Event ایجاد میکند. مدل هوش مصنوعی در Worker فراخوانی میشود.
مزایا:
- پاسخ اولیه API سریعتر است.
- Workerها مستقل مقیاس میشوند.
- اختلال موقت مدل، Threadهای API را برای مدت طولانی اشغال نمیکند.
- Retry در Worker انجام میشود.
- نتیجه موفق و ناموفق Topic جداگانه دارد.
- چند سرویس میتوانند Event نتیجه را مصرف کنند.
- میتوان پردازشهای گذشته را دوباره اجرا کرد.
- مدل Worker بدون تغییر Producer قابل تعویض است.
این معماری برای پردازش فایل، تحلیل دستهای، تولید محتوا، OCR، ترجمه، Embedding، خلاصهسازی و Pipelineهای چندمرحلهای مناسب است.
طراحی Schema مناسب برای Event
Event را فقط به Payload محدود نکنید. Metadata مناسب، Debug و ردیابی را سادهتر میکند.
ساختار پیشنهادی:
{
"event_id": "شناسه یکتای Event",
"event_type": "ai.text.requested",
"event_version": 1,
"occurred_at": "2026-08-06T12:00:00Z",
"trace_id": "شناسه کل جریان",
"causation_id": "شناسه Event قبلی",
"producer": "ai-job-api",
"payload": {
"job_id": "job-123",
"operation": "summarize",
"text": "..."
}
}
نقش فیلدها:
| فیلد | کاربرد |
|---|---|
event_id | تشخیص Event یکتا و Idempotency |
event_type | تعیین نوع Event |
event_version | مدیریت تغییر Schema |
occurred_at | زمان وقوع |
trace_id | اتصال رویدادهای یک Workflow |
causation_id | مشخصکردن Event ایجادکننده |
producer | شناسایی سرویس تولیدکننده |
payload | داده اصلی رویداد |
Schema Evolution
ساختار Event در طول زمان تغییر میکند. حذف یا تغییر ناگهانی یک Field ممکن است Consumerهای قدیمی را خراب کند.
قواعد پیشنهادی:
- Fieldهای جدید را در صورت امکان Optional اضافه کنید.
- معنی Field موجود را تغییر ندهید.
- نسخه Event را ثبت کنید.
- Consumer را نسبت به Fieldهای اضافه مقاوم کنید.
- برای تغییر ناسازگار، نسخه جدید Event بسازید.
- نمونه Eventها را در تست قرارداد نگه دارید.
- Producer و Consumer را مستقل Deploy و آزمایش کنید.
نمونه نسخه جدید:
ai.text.requested.v1
ai.text.requested.v2
در پروژههای بزرگ میتوان از JSON Schema، Avro یا Protobuf و Schema Registry استفاده کرد.
Retry و Dead Letter Topic
Retry نامحدود یک Event خراب میتواند Consumer را متوقف کند. بهتر است خطاها دستهبندی شوند.
خطای موقت
نمونهها:
- Timeout
- پاسخ موقت ۵xx
- محدودیت لحظهای سرویس
- قطعی کوتاه شبکه
برای این خطاها Retry محدود با Backoff مناسب است.
خطای دائمی
نمونهها:
- Payload نامعتبر
- Operation پشتیبانینشده
- متن خالی
- Schema ناسازگار
- شناسه مدل نامعتبر
- دادهای که Validation را رد میکند
این Eventها بهتر است پس از ثبت اطلاعات کافی به Dead Letter Topic منتقل شوند.
Topic پیشنهادی:
ai.text.failed.v1
اطلاعات خطا نباید شامل کلید API یا داده حساس غیرضروری باشد.
Retry Topic برای سیستمهای بزرگتر
در سیستمهای پرترافیک میتوان چند Retry Topic داشت:
ai.text.retry.10s.v1
ai.text.retry.1m.v1
ai.text.retry.10m.v1
ai.text.failed.v1
Event پس از خطا به Topic متناسب منتقل و با تأخیر دوباره پردازش میشود. Kafka بهصورت پایه یک Scheduler عمومی برای Delay Queue نیست؛ بنابراین پیادهسازی Retry زماندار نیازمند Consumer، Timestamp، ابزار تکمیلی یا طراحی جداگانه است.
مقیاسکردن Workerها
فرض کنید Topic سه Partition دارد. میتوانید سه Worker با Group یکسان اجرا کنید:
python worker.py
python worker.py
python worker.py
همه Workerها دارای این Group هستند:
ai-text-workers-v1
Kafka Partitionها را میان آنها توزیع میکند.
اگر تعداد Workerها از Partitionها بیشتر باشد، Workerهای اضافه ظرفیت مصرف مستقیم ندارند. برای افزایش Parallelism ممکن است لازم باشد تعداد Partitionها را افزایش دهید.
افزایش Partition تصمیمی مهم است، زیرا:
- نحوه توزیع Keyها تغییر میکند.
- ترتیب سراسری وجود ندارد.
- تعداد زیاد Partition هزینه مدیریتی دارد.
- کاهش تعداد Partition ساده نیست.
- تغییر Partition میتواند Mapping Key به Partition را تغییر دهد.
Rebalancing چیست؟
وقتی Consumer جدید به Group اضافه میشود، Consumer خارج میشود یا تعداد Partitionها تغییر میکند، Kafka Assignmentها را دوباره توزیع میکند. این فرایند Rebalancing نام دارد.
Rebalance زیاد میتواند باعث توقفهای کوتاه در پردازش شود. علتهای رایج:
- پردازش هر Message بیش از
max.poll.interval.msطول میکشد. - Consumer مرتب Restart میشود.
- Timeoutها درست تنظیم نشدهاند.
- Worker قبل از Poll بعدی کار بسیار طولانی انجام میدهد.
- Deployها همه Consumerها را همزمان قطع میکنند.
برای Jobهای بسیار طولانی، باید Poll Loop و پردازش را با دقت طراحی کنید یا عملیات سنگین را به لایه مناسبتری منتقل کنید.
Backpressure در Pipeline هوش مصنوعی
ممکن است Producer سریعتر از Workerها Event تولید کند. در این حالت Consumer Lag افزایش پیدا میکند.
Consumer Lag تفاوت میان آخرین Offset موجود و Offset پردازششده Consumer است.
افزایش Lag میتواند نشان دهد:
- تعداد Workerها کم است.
- مدل هوش مصنوعی کند پاسخ میدهد.
- Rate Limit مدل محدودکننده است.
- Payloadها بسیار بزرگاند.
- Retry زیاد شده است.
- Partition کافی وجود ندارد.
- Consumer در حال خطاست.
راهکارها:
- افزایش Worker تا سقف Parallelism Topic
- افزایش Partition با ارزیابی دقیق
- Batch Processing برای عملیات مناسب
- استفاده از مدل سریعتر
- محدودکردن طول ورودی
- Route کردن Jobها براساس پیچیدگی
- تنظیم Concurrency با توجه به Rate Limit
- اضافهکردن Admission Control
- اعلام زمان تقریبی انتظار به کاربر
انتخاب مدل هوش مصنوعی در Worker
همه Jobها به یک مدل قدرتمند و گران نیاز ندارند. میتوانید Event را براساس Operation یا پیچیدگی Route کنید:
خلاصهسازی کوتاه → مدل سریع و اقتصادی
تحلیل پیچیده → مدل Reasoning
پردازش تصویر → مدل Vision
تولید کد → مدل Coding
شناسه مدل بهتر است در Configuration نگهداری شود، نه داخل Event ورودی کاربر. اگر کاربر اجازه انتخاب مدل دارد، باید مدل انتخابی با فهرست مجاز Backend اعتبارسنجی شود.
برای بررسی مدلها و قیمتها به صفحه مدلهای درواره مراجعه کنید.
مانیتورینگ Kafka
در Production فقط سلامت Broker کافی نیست. این شاخصها را مانیتور کنید:
شاخصهای Broker
- تعداد Brokerهای فعال
- Offline Partition
- Under-Replicated Partition
- Disk Usage
- Network Throughput
- Request Latency
- CPU و Memory
- زمان Garbage Collection
- نرخ Produce و Fetch
شاخصهای Consumer
- Consumer Lag
- نرخ پردازش
- زمان پردازش هر Event
- تعداد Rebalance
- نرخ خطا
- تعداد Retry
- تعداد Eventهای Dead Letter
- زمان آخرین Commit
- تعداد Consumerهای فعال
شاخصهای AI Worker
- زمان پاسخ مدل
- تعداد Token ورودی و خروجی
- هزینه هر Job
- نرخ پاسخ موفق
- Rate Limit
- Timeout
- تعداد Retry
- تعداد Jobهای ناموفق
- طول Queue
- P50، P95 و P99 Latency
لاگگذاری مناسب
برای هر مرحله این شناسهها را ثبت کنید:
event_id
job_id
trace_id
topic
partition
offset
consumer_group
operation
model
attempt
duration_ms
status
کل متن کاربر، کل پاسخ مدل، کلید API و اطلاعات حساس را بدون نیاز مشخص در Log قرار ندهید.
نمونه Log ساختاریافته:
{
"level": "info",
"event": "ai_job_completed",
"job_id": "job-123",
"trace_id": "job-123",
"topic": "ai.text.requested.v1",
"partition": 1,
"offset": 1259,
"duration_ms": 2840,
"status": "completed"
}
تستکردن Kafka Pipeline
تست Producer
بررسی کنید:
- Event در Topic درست منتشر میشود.
- Key مناسب دارد.
- Schema معتبر است.
- در نبود Kafka پاسخ کنترلشده برمیگردد.
- Payload بزرگ یا نامعتبر رد میشود.
تست Consumer
بررسی کنید:
- Event معتبر پردازش میشود.
- Event نامعتبر به مسیر خطا میرود.
- پس از خطای موقت Retry انجام میشود.
- Offset زودتر از موعد Commit نمیشود.
- پردازش تکراری نتیجه نامعتبر ایجاد نمیکند.
- Shutdown باعث ازدسترفتن وضعیت نمیشود.
تست Integration
یک Kafka واقعی یا Container آزمایشی اجرا کنید و مسیر کامل را تست کنید:
HTTP Request
→ Request Topic
→ Worker
→ API Mock
→ Completed Topic
→ Result Consumer
تست Load
موارد زیر را اندازه بگیرید:
- تعداد Event در ثانیه
- Producer Latency
- Consumer Lag
- زمان تکمیل Job
- رفتار هنگام کندشدن مدل
- رفتار هنگام افزایش Error Rate
- سرعت بازیابی پس از Restart Worker
خطاهای رایج Kafka
خطای Connection refused
علتهای احتمالی:
- Kafka اجرا نشده است.
- پورت ۹۰۹۲ در دسترس نیست.
bootstrap.serversاشتباه است.- برنامه داخل Docker از
localhostاستفاده میکند.
بررسی:
docker compose ps
docker compose logs kafka
اگر برنامه Python داخل همان Docker Compose باشد، آدرس معمولاً باید نام Service باشد:
kafka:9092
اما برنامهای که روی Host اجرا میشود از این آدرس استفاده میکند:
localhost:9092
خطای Topic not found
فهرست Topicها را ببینید:
docker exec darvareh-kafka \
/opt/kafka/bin/kafka-topics.sh \
--list \
--bootstrap-server localhost:9092
نام Topic در .env باید دقیقاً با Topic ساختهشده یکسان باشد.
Consumer پیامی دریافت نمیکند
موارد زیر را بررسی کنید:
- Topic درست است.
- Producer واقعاً Event فرستاده است.
- Consumer Group قبلاً Offset را Commit نکرده است.
- مقدار
auto.offset.resetفقط برای Group بدون Offset قبلی اثر دارد. - Consumer دیگری در همان Group Partition را گرفته است.
- Event در Partition دیگری قرار دارد.
- Worker دچار Rebalance مداوم نشده است.
برای تست از یک group.id جدید استفاده کنید.
Event دوباره پردازش میشود
این رفتار در مدل At-least-once امکانپذیر است. راهحل اصلی طراحی Idempotent است، نه فرض اینکه Event هرگز تکرار نمیشود.
موارد زیر را بررسی کنید:
- Offset بعد از Side Effect Commit میشود.
event_idثبت شده است.- Result Storage Unique Constraint دارد.
- API خارجی Idempotency Key میپذیرد.
- Worker میان انتشار نتیجه و Commit متوقف نشده است.
Consumer Lag زیاد میشود
دلایل احتمالی:
- Worker کند است.
- تعداد Partition کم است.
- Rate Limit مدل پایین است.
- Eventها بزرگاند.
- خطا و Retry زیاد شده است.
- Consumer مرتب Rebalance میشود.
- پردازش Blocking طولانی است.
ابتدا Bottleneck را اندازهگیری کنید و بعد تعداد Worker یا Partition را افزایش دهید.
پیامها خارج از ترتیب دیده میشوند
Kafka فقط داخل هر Partition ترتیب را حفظ میکند. بررسی کنید Eventهای مرتبط با Key یکسان منتشر شدهاند.
اندازه Message بیش از حد مجاز است
Kafka برای انتقال فایلهای بزرگ طراحی نشده است. بهجای قرار دادن فایل کامل در Event، فایل را در Object Storage ذخیره و فقط Reference آن را منتشر کنید:
{
"file_id": "file-123",
"object_key": "uploads/file-123.pdf"
}
افزایش بیرویه محدودیت Message Size میتواند مصرف حافظه و شبکه را بیشتر کند.
چکلیست Production
پیش از استفاده Production این موارد را بررسی کنید:
- Topicها نامگذاری و نسخهبندی شدهاند.
- تعداد Partition براساس Throughput انتخاب شده است.
- Replication Factor متناسب با Cluster است.
- Producer از
acks=allاستفاده میکند. - Producer Idempotence در صورت نیاز فعال است.
- Eventها شناسه یکتا دارند.
- Partition Key آگاهانه انتخاب شده است.
- Schema و نسخه Event مشخص است.
- Consumerها Idempotent هستند.
- Offset پس از Side Effect موفق Commit میشود.
- Retry محدود و قابلاندازهگیری است.
- Dead Letter Topic وجود دارد.
- Event خراب پردازش را برای همیشه متوقف نمیکند.
- Consumer Lag مانیتور میشود.
- Rebalance و خطاها Alert دارند.
- Shutdown برنامه Graceful است.
- داده حساس غیرضروری وارد Event نمیشود.
- Messageهای بزرگ به Object Storage منتقل میشوند.
- تست Load انجام شده است.
- سیاست Retention مشخص است.
- برنامه بازیابی بعد از قطعی آزمایش شده است.
- هزینه و زمان پردازش هوش مصنوعی اندازهگیری میشود.
- کلید API فقط در Backend و متغیر محیطی نگهداری میشود.
- مدل انتخابی از Configuration معتبر خوانده میشود.
پرسشهای متداول
Kafka چیست؟
Apache Kafka یک پلتفرم توزیعشده Event Streaming است که برای نوشتن، ذخیره، خواندن و پردازش جریان رویدادها استفاده میشود.
آیا Kafka یک Message Broker است؟
Kafka میتواند نقش Message Broker را ایفا کند، اما معماری آن بیشتر بر Distributed Log و Event Streaming متمرکز است. Retention و Replay از تفاوتهای مهم آن با بسیاری از Queueهای سنتی هستند.
Topic در Kafka چیست؟
Topic یک جریان نامگذاریشده از Eventهاست. Producer در Topic مینویسد و Consumer از آن میخواند.
Partition در Kafka چیست؟
Partition بخشی از Topic و یک Log مرتب است. Partitionها امکان پردازش موازی و توزیع داده را فراهم میکنند.
Offset چیست؟
Offset موقعیت یک Event داخل Partition است. Consumer از Offset برای دنبالکردن میزان پیشرفت پردازش استفاده میکند.
Consumer Group چیست؟
Consumer Group مجموعه Consumerهایی است که Partitionهای Topic را میان خود تقسیم میکنند. Groupهای متفاوت میتوانند یک Topic را مستقل مصرف کنند.
آیا Kafka ترتیب Messageها را تضمین میکند؟
Kafka ترتیب را داخل هر Partition تضمین میکند، نه میان تمام Partitionهای Topic. Eventهای مرتبط باید با Key یکسان منتشر شوند.
Kafka بهتر است یا RabbitMQ؟
هیچکدام همیشه بهتر نیستند. Kafka برای Event Streaming، Replay، جریانهای پرحجم و چند Consumer مستقل مناسب است. RabbitMQ برای Routing پیام و Task Queueهای کلاسیک معمولاً سادهتر است.
آیا Kafka برای پروژه کوچک مناسب است؟
اگر پروژه فقط چند Job ساده دارد، Kafka ممکن است پیچیدگی غیرضروری ایجاد کند. Redis Queue، Celery یا RabbitMQ میتوانند انتخاب سادهتری باشند.
آیا Kafka به ZooKeeper نیاز دارد؟
نسخههای جدید Kafka از KRaft استفاده میکنند و به ZooKeeper نیاز ندارند. آموزش این مقاله بر Kafka 4.3.1 مبتنی است.
آیا میتوان Kafka را با Python استفاده کرد؟
بله. کتابخانههایی مانند confluent-kafka امکان ساخت Producer، Consumer و Admin Client را در Python فراهم میکنند.
آیا میتوان Kafka را به FastAPI متصل کرد؟
بله. FastAPI میتواند درخواست را دریافت و با Producer در Kafka منتشر کند. Workerهای جداگانه نیز Eventها را پردازش میکنند.
آیا Kafka برای پردازش هوش مصنوعی مناسب است؟
برای پردازشهای غیرهمزمان، حجیم یا چندمرحلهای مناسب است. Kafka میتواند درخواستها را میان Workerها توزیع و جریان نتایج را برای سرویسهای دیگر منتشر کند.
آیا Kafka مانع پردازش تکراری میشود؟
نه در همه شرایط. در پردازش At-least-once امکان تکرار وجود دارد. Consumer باید براساس event_id یا Idempotency Key طراحی شود.
هزینه استفاده از API درواره در Worker چقدر است؟
هزینه به مدل، تعداد Tokenهای ورودی و خروجی و تعداد Jobها بستگی دارد. قیمت بهروز مدلها در صفحه مدلهای درواره قابل مشاهده است.
جمعبندی
Apache Kafka یکی از ابزارهای اصلی برای ساخت معماری Event-Driven و پردازش جریان رویدادهاست. مفاهیمی مانند Topic، Partition، Offset، Producer، Consumer و Consumer Group پایههای کار با Kafka را تشکیل میدهند.
در پروژه عملی این مقاله، Kafka را با Docker راهاندازی کردیم، سه Topic برای درخواست، نتیجه و خطا ساختیم و یک Pipeline کامل با Python و FastAPI پیادهسازی کردیم. FastAPI درخواست کاربر را به Event تبدیل میکند، Worker رویداد را از Kafka میخواند و متن را از طریق API درواره پردازش میکند. نتیجه نیز در Topic جداگانه قرار میگیرد.
این معماری برای خلاصهسازی، ترجمه، تحلیل متن، پردازش فایل، تولید محتوا، Embedding و Workflowهای غیرهمزمان هوش مصنوعی قابل توسعه است. بااینحال، Kafka زمانی انتخاب مناسبی است که نیاز واقعی به Event Streaming، Replay، توزیع پردازش یا جداسازی سرویسها وجود داشته باشد.
برای شروع استفاده از مدلهای مختلف هوش مصنوعی، در درواره ثبتنام کنید، کلید API بسازید و مدل مناسب Workerهای خود را از صفحه مدلها و قیمتها انتخاب کنید.
منابع پیشنهادی
- مستندات رسمی Apache Kafka
- راهنمای شروع سریع Apache Kafka
- مستندات Kafka Streams
- مستندات Python Client برای Kafka
- مستندات FastAPI
- مستندات Docker Compose
مقالات مرتبط
- پردازش غیرهمزمان API هوش مصنوعی با Job Queue، Worker و Webhook
- ساخت API هوش مصنوعی آماده Production
- مانیتورینگ و Observability در سیستمهای هوش مصنوعی
- آموزش اتصال API هوش مصنوعی به اپلیکیشن
- تست Load و Performance API با k6 و هوش مصنوعی
- تحلیل بازخورد مشتریان با هوش مصنوعی
برای مطالعه شرایط استفاده و محدودیتهای مسئولیت، صفحه «سلب مسئولیت» را مشاهده کنید.