Kafka چیست؟ آموزش کامل Apache Kafka، Event Streaming و اتصال به Python و FastAPI

Apache Kafka پلتفرمی برای انتقال و پردازش جریان رویدادهاست. در این آموزش، Kafka را با Docker راه‌اندازی می‌کنید و یک سیستم Event-Driven عملی با Python، FastAPI و API هوش مصنوعی درواره می‌سازید.

Share
Kafka چیست؟ آموزش کامل Apache Kafka، Event Streaming و اتصال به Python و FastAPI

وقتی یک برنامه کوچک است، معمولاً همه عملیات در همان درخواست 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 چه مشکلی را حل می‌کند؟

فرض کنید کاربر متنی را برای تحلیل هوش مصنوعی ارسال می‌کند. انجام همه مراحل در یک درخواست ممکن است شامل این عملیات باشد:

  1. اعتبارسنجی متن
  2. ذخیره درخواست
  3. ارسال به مدل هوش مصنوعی
  4. دریافت خروجی
  5. ذخیره نتیجه
  6. ارسال Notification
  7. ثبت آمار مصرف
  8. به‌روزرسانی داشبورد

اگر همه این مراحل در یک 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 KafkaQueue سنتی
مدل اصلیEvent Log توزیع‌شدهصف پیام
نگهداری پس از مصرفبراساس Retentionمعمولاً حذف پس از Ack
Replayقابلیت اصلیمعمولاً محدودتر
ترتیبداخل هر Partitionوابسته به Queue
پردازش موازیبا Partition و Consumer Groupبا Workerهای صف
چند Consumer مستقلبا Groupهای مختلفبسته به ابزار و Exchange
کاربرد مناسبEvent Streaming، Integration، AnalyticsTask Queue و Job Processing
مقیاس جریان دادهبسیار بالاوابسته به محصول

انتخاب Kafka فقط به‌دلیل محبوبیت آن تصمیم مناسبی نیست. اگر صرفاً چند Job پس‌زمینه ساده دارید، ابزارهایی مانند Redis Queue، RabbitMQ یا Celery ممکن است ساده‌تر باشند.

مقایسه Kafka، RabbitMQ و Redis Streams

معیارKafkaRabbitMQRedis Streams
کاربرد اصلیEvent StreamingMessage BrokerStream سبک در 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

در این پروژه سه برنامه داریم:

  1. producer_api.py درخواست HTTP را دریافت و Event را منتشر می‌کند.
  2. worker.py Event را مصرف و متن را از طریق API درواره پردازش می‌کند.
  3. 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های خود را از صفحه مدل‌ها و قیمت‌ها انتخاب کنید.

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

مقالات مرتبط

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

Read more

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

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

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

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

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

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