~/icsd.ir — bash
SYSTEM_ONLINE

پیام‌رسانی Async پیشرفته

در فصل ۳ مقدمه‌ای از پیام‌رسانی async دیدیم. در این فصل عمیق‌تر می‌رویم: الگوهای پیشرفته، انتخاب RabbitMQ vs Kafka، patterns مختلف، و پیاده‌سازی production-grade.

۶.۱ مقدمه

در فصل ۳ مقدمه‌ای از پیام‌رسانی async دیدیم. در این فصل عمیق‌تر می‌رویم: الگوهای پیشرفته، انتخاب RabbitMQ vs Kafka، patterns مختلف، و پیاده‌سازی production-grade.


هدف این فصل: تسلط بر RabbitMQ و Kafka، آشنایی با Redis Pub/Sub، الگوهای Work Queue، RPC، Publish/Subscribe، Routing، Topic، Headers و Schema Registry.

۶.۲ الگوهای پیام‌رسانی

۱. Point-to-Point (Queue)

یک پیام فقط توسط یک consumer دریافت می‌شود. اگر چند consumer داشته باشید، پیام‌ها بین آن‌ها distribute می‌شوند.


[Producer] ──► [Queue] ──┬──► [Consumer 1] (پیام‌های فرد)
                          ├──► [Consumer 2] (پیام‌های زوج)
                          └──► [Consumer 3] (...)
    

کاربرد: Work Queue — پردازش task ها بین چند worker.

۲. Publish/Subscribe (Topic)

یک پیام به همه subscriber ها می‌رود.


[Producer] ──► [Topic] ──┬──► [Subscriber 1] (همه پیام‌ها)
                          ├──► [Subscriber 2] (همه پیام‌ها)
                          └──► [Subscriber 3] (همه پیام‌ها)
    

کاربرد: Event broadcasting — اطلاع‌رسانی به چند سرویس همزمان.

۳. Request/Reply (RPC)

Producer منتظر پاسخ می‌ماند. مثل RPC ولی روی message queue.


[Client] ──(req)──► [Request Queue]
              ▲                 │
              │                 ▼
              │           [RPC Server]
              │                 │
              │           [Reply Queue]
              └◄──(rsp)─────────┘
    

۴. Routing

پیام بر اساس routing key به queue های مختلف می‌رود.


[Producer] ──► [Direct Exchange] ──┬─(error)──► [Error Queue]
                                    ├─(warn)───► [Warning Queue]
                                    └─(info)───► [Info Queue]
    

۵. Competing Consumers

چند consumer با هم پیام‌ها را پردازش می‌کنند برای parallelism.

۶. Priority Queue

پیام‌های مهم‌تر اول پردازش می‌شوند.

۷. Delayed Messages

پیام پس از مدت زمان مشخص پردازش می‌شود.

۶.۳ RabbitMQ — عمیق‌تر

اجرای کامل RabbitMQ


# docker-compose.yml
version: "3.8"

services:
  rabbitmq:
    image: rabbitmq:3.12-management
    ports:
      - "5672:5672"   # AMQP
      - "15672:15672" # Management UI
    environment:
      RABBITMQ_DEFAULT_USER: admin
      RABBITMQ_DEFAULT_PASS: admin
      RABBITMQ_DEFAULT_VHOST: /
    volumes:
      - rabbitmq-data:/var/lib/rabbitmq
      - ./rabbitmq.conf:/etc/rabbitmq/rabbitmq.conf
      - ./definitions.json:/etc/rabbitmq/definitions.json
    networks:
      - microservices

volumes:
  rabbitmq-data:
    

تنظیمات اضافی


# rabbitmq.conf
default_user = admin
default_pass = admin
default_vhost = /

# Memory و disk
vm_memory_high_watermark.relative = 0.6
disk_free_limit.relative = 2.0

# Performance
heartbeat = 60
frame_max = 131072
channel_max = 2047

# Plugins
management.tcp.port = 15672

# اعلامیه‌ها
loopback_users.guest = false

# Cluster (در صورت نیاز)
# cluster_formation.peer_discovery_backend = classic_config
    

پیاده‌سازی Production-Grade Producer


# rabbitmq/publisher.py
import pika
import json
import logging
from typing import Optional, Dict, Any
from contextlib import contextmanager

logger = logging.getLogger(__name__)

class RabbitMQPublisher:
    """Publisher با connection pooling و retry"""
    
    def __init__(
        self,
        host: str = "rabbitmq",
        port: int = 5672,
        username: str = "admin",
        password: str = "admin",
        vhost: str = "/",
    ):
        self.params = pika.ConnectionParameters(
            host=host,
            port=port,
            virtual_host=vhost,
            credentials=pika.PlainCredentials(username, password),
            heartbeat=600,
            blocked_connection_timeout=300,
            connection_attempts=3,
            retry_delay=5,
        )
        self._connection: Optional[pika.BlockingConnection] = None
        self._channel: Optional[pika.adapters.blocking_connection.BlockingChannel] = None
    
    def _ensure_connection(self):
        """اطمینان از connection سالم"""
        if self._connection is None or self._connection.is_closed:
            logger.info("Creating new RabbitMQ connection")
            self._connection = pika.BlockingConnection(self.params)
            self._channel = self._connection.channel()
            self._channel.confirm_delivery()  # publisher confirms
    
    def declare_exchange(
        self,
        name: str,
        exchange_type: str = "topic",
        durable: bool = True
    ):
        """تعریف exchange"""
        self._ensure_connection()
        self._channel.exchange_declare(
            exchange=name,
            exchange_type=exchange_type,
            durable=durable,
        )
    
    def publish(
        self,
        exchange: str,
        routing_key: str,
        message: Dict[str, Any],
        headers: Optional[Dict] = None,
        priority: int = 0,
        expiration: Optional[int] = None,  # ms
        persistent: bool = True,
    ) -> bool:
        """انتشار پیام"""
        self._ensure_connection()
        
        properties = pika.BasicProperties(
            content_type="application/json",
            delivery_mode=2 if persistent else 1,
            priority=priority,
            headers=headers or {},
            expiration=str(expiration) if expiration else None,
            message_id=message.get("id"),
            timestamp=int(time.time()),
        )
        
        try:
            self._channel.basic_publish(
                exchange=exchange,
                routing_key=routing_key,
                body=json.dumps(message),
                properties=properties,
                mandatory=True,  # خطا اگر unrouted
            )
            logger.info(f"Published to {exchange}/{routing_key}")
            return True
        except pika.exceptions.UnroutableError:
            logger.error(f"Message unrouted: {exchange}/{routing_key}")
            return False
        except pika.exceptions.AMQPConnectionError:
            logger.error("Connection lost, retrying...")
            self._connection = None  # force reconnect
            raise
    
    def close(self):
        if self._connection and not self._connection.is_closed:
            self._connection.close()

# Singleton instance
publisher = RabbitMQPublisher()

# استفاده
def publish_order_event(order_id: int, event_type: str, data: dict):
    publisher.declare_exchange("orders", "topic", durable=True)
    
    message = {
        "id": f"{order_id}-{event_type}-{int(time.time())}",
        "event_type": event_type,
        "order_id": order_id,
        "data": data,
        "timestamp": time.time(),
    }
    
    publisher.publish(
        exchange="orders",
        routing_key=f"order.{event_type}",
        message=message,
        headers={
            "x-source-service": "order-service",
            "x-trace-id": get_trace_id(),
        }
    )
    

Production-Grade Consumer


# rabbitmq/consumer.py
import pika
import json
import logging
import time
from typing import Callable

logger = logging.getLogger(__name__)

class RabbitMQConsumer:
    def __init__(
        self,
        queue_name: str,
        host: str = "rabbitmq",
        prefetch_count: int = 10,
        max_retries: int = 3,
    ):
        self.queue_name = queue_name
        self.host = host
        self.prefetch_count = prefetch_count
        self.max_retries = max_retries
        self._handlers: dict = {}
    
    def register_handler(self, routing_key: str, handler: Callable):
        """ثبت handler برای routing key خاص"""
        self._handlers[routing_key] = handler
    
    def setup_queue(self, channel, exchange: str, routing_keys: list):
        """تنظیم queue با DLQ"""
        # Dead Letter Exchange
        dlx_name = f"{exchange}.dlx"
        dlq_name = f"{self.queue_name}.dlq"
        
        channel.exchange_declare(exchange=dlx_name, exchange_type="direct", durable=True)
        channel.queue_declare(queue=dlq_name, durable=True)
        channel.queue_bind(exchange=dlx_name, queue=dlq_name, routing_key=dlq_name)
        
        # Main queue با DLX
        channel.queue_declare(
            queue=self.queue_name,
            durable=True,
            arguments={
                "x-dead-letter-exchange": dlx_name,
                "x-dead-letter-routing-key": dlq_name,
                "x-max-priority": 10,  # priority queue
            }
        )
        
        # Bind به routing keys
        for rk in routing_keys:
            channel.queue_bind(
                exchange=exchange,
                queue=self.queue_name,
                routing_key=rk,
            )
    
    def _process_message(self, ch, method, properties, body):
        """پردازش یک پیام"""
        retry_count = 0
        if properties.headers:
            retry_count = properties.headers.get("x-retry-count", 0)
        
        try:
            message = json.loads(body)
            routing_key = method.routing_key
            
            # پیدا کردن handler
            handler = self._handlers.get(routing_key)
            if not handler:
                # wildcard match
                for key, h in self._handlers.items():
                    if self._matches(key, routing_key):
                        handler = h
                        break
            
            if not handler:
                logger.warning(f"No handler for {routing_key}")
                ch.basic_ack(delivery_tag=method.delivery_tag)
                return
            
            # پردازش
            handler(message)
            
            # موفق
            ch.basic_ack(delivery_tag=method.delivery_tag)
            logger.info(f"Processed {message.get('id')} [{routing_key}]")
            
        except Exception as e:
            logger.error(f"Error processing message: {e}", exc_info=True)
            
            if retry_count < self.max_retries:
                # republish با retry count
                new_headers = (properties.headers or {}).copy()
                new_headers["x-retry-count"] = retry_count + 1
                new_headers["x-last-error"] = str(e)[:200]
                
                # Exponential backoff
                delay = (2 ** retry_count) * 1000  # ms
                
                ch.basic_publish(
                    exchange="",
                    routing_key=method.routing_key,
                    body=body,
                    properties=pika.BasicProperties(
                        headers=new_headers,
                        expiration=str(delay),
                    )
                )
                ch.basic_ack(delivery_tag=method.delivery_tag)
                logger.info(f"Scheduled retry {retry_count + 1}/{self.max_retries}")
            else:
                # به DLQ
                ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
                logger.error(f"Max retries exceeded, sent to DLQ")
    
    def _matches(self, pattern: str, key: str) -> bool:
        """تطابق routing key با pattern"""
        pattern_parts = pattern.split(".")
        key_parts = key.split(".")
        
        for p, k in zip(pattern_parts, key_parts):
            if p == "#":
                return True
            if p != "*" and p != k:
                return False
        
        return len(pattern_parts) == len(key_parts)
    
    def consume(self, exchange: str, routing_keys: list):
        """شروع مصرف"""
        params = pika.ConnectionParameters(
            host=self.host,
            credentials=pika.PlainCredentials("admin", "admin"),
            heartbeat=600,
        )
        connection = pika.BlockingConnection(params)
        channel = connection.channel()
        
        # تنظیم queue
        self.setup_queue(channel, exchange, routing_keys)
        
        # QoS
        channel.basic_qos(prefetch_count=self.prefetch_count)
        
        # شروع مصرف
        channel.basic_consume(
            queue=self.queue_name,
            on_message_callback=self._process_message,
            auto_ack=False,
        )
        
        logger.info(f"Consuming from {self.queue_name}...")
        try:
            channel.start_consuming()
        except KeyboardInterrupt:
            channel.stop_consuming()
            connection.close()

# استفاده
def handle_order_placed(message):
    order_id = message["order_id"]
    print(f"Sending confirmation email for order {order_id}")
    # ...

def handle_order_cancelled(message):
    order_id = message["order_id"]
    print(f"Refunding order {order_id}")
    # ...

consumer = RabbitMQConsumer("notification_orders_queue")
consumer.register_handler("order.placed", handle_order_placed)
consumer.register_handler("order.cancelled", handle_order_cancelled)
consumer.consume("orders", ["order.*"])
    

۶.۴ Apache Kafka — عمیق‌تر

اجرای Kafka Cluster


# docker-compose.yml
version: "3.8"

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    volumes:
      - zk-data:/var/lib/zookeeper/data
      - zk-logs:/var/lib/zookeeper/log

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
      KAFKA_LOG_RETENTION_HOURS: 168  # ۷ روز
      KAFKA_LOG_RETENTION_BYTES: 1073741824  # 1GB
      KAFKA_NUM_PARTITIONS: 6
      KAFKA_DEFAULT_REPLICATION_FACTOR: 1
    volumes:
      - kafka-data:/var/lib/kafka/data

  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092
      KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper:2181
    depends_on:
      - kafka

volumes:
  zk-data:
  zk-logs:
  kafka-data:
    

ایجاد Topic


# Topic با ۶ partition و replication 3
docker exec kafka kafka-topics --create 
    --topic order-events 
    --bootstrap-server kafka:9092 
    --partitions 6 
    --replication-factor 1 
    --config retention.ms=604800000 
    --config compression.type=snappy

# لیست topic ها
docker exec kafka kafka-topics --list --bootstrap-server kafka:9092

# جزئیات topic
docker exec kafka kafka-topics --describe 
    --topic order-events 
    --bootstrap-server kafka:9092
    

Producer پیشرفته


# kafka/producer.py
from kafka import KafkaProducer
from kafka.errors import KafkaError
import json
import logging

logger = logging.getLogger(__name__)

class KafkaPublisher:
    def __init__(
        self,
        bootstrap_servers: list = ["kafka:9092"],
        client_id: str = "order-service",
    ):
        self.producer = KafkaProducer(
            bootstrap_servers=bootstrap_servers,
            client_id=client_id,
            
            # Serialization
            value_serializer=lambda v: json.dumps(v).encode("utf-8"),
            key_serializer=lambda k: str(k).encode("utf-8") if k else None,
            
            # Reliability
            acks="all",  # منتظر تأیید همه replica ها
            retries=5,
            max_in_flight_requests_per_connection=5,
            enable_idempotence=True,  # exactly-once
            
            # Performance
            batch_size=16384,  # 16KB
            linger_ms=10,  # batch تا ۱۰ms
            compression_type="snappy",
            buffer_memory=33554432,  # 32MB
        )
    
    def publish(
        self,
        topic: str,
        event: dict,
        key: str = None,
        headers: list = None,
    ):
        """انتشار پیام"""
        try:
            future = self.producer.send(
                topic=topic,
                value=event,
                key=key,
                headers=headers,
            )
            
            # async callback
            future.add_callback(self._on_success, topic=topic)
            future.add_errback(self._on_error, topic=topic)
            
            return future
        except KafkaError as e:
            logger.error(f"Failed to publish: {e}")
            raise
    
    def publish_sync(self, topic: str, event: dict, key: str = None) -> dict:
        """انتشار همزمان (blocking)"""
        future = self.publish(topic, event, key)
        try:
            metadata = future.get(timeout=10)
            return {
                "topic": metadata.topic,
                "partition": metadata.partition,
                "offset": metadata.offset,
            }
        except KafkaError as e:
            logger.error(f"Sync publish failed: {e}")
            raise
    
    def _on_success(self, metadata, topic):
        logger.debug(
            f"Sent to {topic} partition={metadata.partition} "
            f"offset={metadata.offset}"
        )
    
    def _on_error(self, exception, topic):
        logger.error(f"Error sending to {topic}: {exception}")
    
    def flush(self):
        self.producer.flush(timeout=10)
    
    def close(self):
        self.producer.close(timeout=10)

# استفاده
publisher = KafkaPublisher()

# انتشار event
publisher.publish(
    topic="order-events",
    key=str(order_id),  # برای partitioning
    event={
        "event_type": "order_placed",
        "order_id": order_id,
        "user_id": user_id,
        "total": total,
        "timestamp": time.time(),
    },
    headers=[
        ("source", b"order-service"),
        ("version", b"1.0"),
    ]
)

publisher.flush()
    

Consumer Group پیشرفته


# kafka/consumer.py
from kafka import KafkaConsumer, TopicPartition
import json
import logging
from typing import Callable

logger = logging.getLogger(__name__)

class KafkaConsumerWrapper:
    def __init__(
        self,
        topics: list,
        group_id: str,
        bootstrap_servers: list = ["kafka:9092"],
        auto_offset_reset: str = "earliest",
    ):
        self.consumer = KafkaConsumer(
            *topics,
            bootstrap_servers=bootstrap_servers,
            group_id=group_id,
            
            # Deserialization
            value_deserializer=lambda v: json.loads(v.decode("utf-8")),
            key_deserializer=lambda k: k.decode("utf-8") if k else None,
            
            # Offset
            auto_offset_reset=auto_offset_reset,  # earliest یا latest
            enable_auto_commit=False,  # manual commit برای دقت
            
            # Performance
            max_poll_records=500,
            max_poll_interval_ms=300000,  # 5 min
            session_timeout_ms=10000,
            heartbeat_interval_ms=3000,
            
            # Fetching
            fetch_max_bytes=52428800,  # 50MB
            max_partition_fetch_bytes=1048576,  # 1MB per partition
        )
        self._handlers: dict = {}
    
    def register_handler(self, event_type: str, handler: Callable):
        self._handlers[event_type] = handler
    
    def consume(self):
        """consumption loop"""
        try:
            for message in self.consumer:
                self._process_message(message)
                
                # commit دستی
                self.consumer.commit()
        except KeyboardInterrupt:
            logger.info("Stopping consumer")
        finally:
            self.consumer.close()
    
    def _process_message(self, message):
        """پردازش یک پیام"""
        try:
            event = message.value
            event_type = event.get("event_type")
            
            handler = self._handlers.get(event_type)
            if not handler:
                logger.warning(f"No handler for {event_type}")
                return
            
            logger.info(
                f"Processing {event_type} from "
                f"partition={message.partition} offset={message.offset}"
            )
            
            handler(event)
            
        except Exception as e:
            logger.error(f"Failed to process message: {e}", exc_info=True)
            # در Kafka، اگر commit نکنیم، دوباره می‌خوانیم
            # برای جلوگیری از infinite loop:
            self._send_to_dlq(message, str(e))
    
    def _send_to_dlq(self, message, error):
        """فرستادن به Dead Letter Topic"""
        from .producer import KafkaPublisher
        publisher = KafkaPublisher(client_id="dlq-publisher")
        publisher.publish(
            topic=f"{message.topic}.dlq",
            key=message.key,
            event={
                "original_value": message.value,
                "error": error,
                "topic": message.topic,
                "partition": message.partition,
                "offset": message.offset,
                "failed_at": time.time(),
            }
        )
        publisher.flush()

# استفاده
def handle_order_placed(event):
    print(f"Email for order {event['order_id']}")

def handle_order_cancelled(event):
    print(f"Refund for order {event['order_id']}")

consumer = KafkaConsumerWrapper(
    topics=["order-events"],
    group_id="notification-service",
)
consumer.register_handler("order_placed", handle_order_placed)
consumer.register_handler("order_cancelled", handle_order_cancelled)
consumer.consume()
    

۶.۵ Redis Pub/Sub

Redis یک سیستم pub/sub ساده و سریع دارد. برای real-time messaging سبک عالی است (مثل notifications، WebSocket fanout).


# redis_pubsub.py
import redis.asyncio as aioredis
import json
import asyncio

class RedisPubSub:
    def __init__(self, url: str = "redis://redis:6379"):
        self.redis = aioredis.from_url(url)
    
    async def publish(self, channel: str, message: dict):
        await self.redis.publish(channel, json.dumps(message))
    
    async def subscribe(self, channel: str, handler):
        pubsub = self.redis.pubsub()
        await pubsub.subscribe(channel)
        
        async for message in pubsub.listen():
            if message["type"] == "message":
                data = json.loads(message["data"])
                await handler(data)

# استفاده
async def on_notification(data):
    print(f"Got: {data}")

ps = RedisPubSub()
await ps.subscribe("user:123:notifications", on_notification)

# در سرویس دیگر:
await ps.publish("user:123:notifications", {
    "type": "order_shipped",
    "order_id": 456
})
    

تفاوت با RabbitMQ/Kafka

  • Redis: پیام‌ها persist نمی‌شوند — اگر subscriber آنلاین نباشد، از دست می‌رود
  • Redis: ساده‌ترین، سریع‌ترین
  • Redis Streams: نسخه persistent (شبیه Kafka)

Redis Streams (مدرن‌تر)


# Redis Streams - persistent
async def add_to_stream(redis, stream: str, data: dict):
    await redis.xadd(stream, data)

async def consume_stream(redis, stream: str, group: str, consumer: str):
    # ایجاد consumer group
    try:
        await redis.xgroup_create(stream, group, id="0", mkstream=True)
    except Exception:
        pass  # group already exists
    
    while True:
        messages = await redis.xreadgroup(
            group, consumer,
            {stream: ">"},
            count=10,
            block=5000
        )
        
        for stream_name, msgs in messages:
            for msg_id, fields in msgs:
                try:
                    process(fields)
                    await redis.xack(stream, group, msg_id)
                except Exception as e:
                    print(f"Error: {e}")
    

۶.۶ RabbitMQ vs Kafka — انتخاب درست

قانون ساده انتخاب

  • RabbitMQ: اگر پیام «task» است (پردازش شود و بسوزد)
  • Kafka: اگر پیام «event» است (تاریخچه مهم است)

سناریوهای کاربردی

سناریو توصیه چرا
ارسال ایمیل، SMS RabbitMQ + Celery Task queue ساده، کافی است
RPC over messaging RabbitMQ request/reply pattern مستقیم
Event sourcing Kafka events باید persist و قابل replay باشند
Real-time analytics Kafka throughput بالا، Kafka Streams
Log aggregation Kafka retention طولانی
Microservice events ساده RabbitMQ سادگی، routing انعطاف‌پذیر
Cross-service events پایدار Kafka چند consumer، replay
Real-time chat Redis Pub/Sub سریع، سبک
Multi-region replication Kafka MirrorMaker
پروژه startup کوچک RabbitMQ یا Redis سادگی، سرعت setup

۶.۷ Schema Registry و Schema Evolution

یکی از چالش‌های مهم پیام‌رسانی: «چگونه schema را تغییر دهیم بدون breaking consumer ها؟»

قوانین Schema Evolution

  • Backward Compatible: consumer قدیمی می‌تواند پیام جدید را بخواند
  • Forward Compatible: consumer جدید می‌تواند پیام قدیمی را بخواند
  • Full Compatible: هر دو

قوانین برای حفظ Compatibility

  • ✅ اضافه کردن field optional (با default)
  • ✅ حذف field optional
  • ❌ حذف field required
  • ❌ تغییر type field
  • ❌ rename کردن field (به جای آن، اضافه کن و قدیمی را deprecate کن)

استفاده از Confluent Schema Registry


# docker-compose.yml
schema-registry:
  image: confluentinc/cp-schema-registry:7.5.0
  ports:
    - "8081:8081"
  environment:
    SCHEMA_REGISTRY_HOST_NAME: schema-registry
    SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: kafka:9092
    SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081
  depends_on:
    - kafka
    

# با Avro
from confluent_kafka import avro
from confluent_kafka.avro import AvroProducer

value_schema_str = """
{
   "namespace": "shop.events",
   "name": "OrderPlaced",
   "type": "record",
   "fields": [
       {"name": "order_id", "type": "int"},
       {"name": "user_id", "type": "int"},
       {"name": "total", "type": "double"},
       {"name": "timestamp", "type": "long"},
       {"name": "currency", "type": "string", "default": "IRR"}
   ]
}
"""

value_schema = avro.loads(value_schema_str)

producer = AvroProducer({
    "bootstrap.servers": "kafka:9092",
    "schema.registry.url": "http://schema-registry:8081"
}, default_value_schema=value_schema)

producer.produce(
    topic="order-events",
    value={
        "order_id": 123,
        "user_id": 5,
        "total": 1500.50,
        "timestamp": int(time.time() * 1000),
    }
)
producer.flush()
    

۶.۸ پیاده‌سازی الگوها

۱. Outbox Pattern

راه‌حل کلاسیک برای dual-write problem: «چطور هم DB را آپدیت کنم و هم event منتشر کنم؟»


# مشکل: dual write
async def create_order_BAD(data):
    order = await db.create_order(data)
    await rabbitmq.publish("order.placed", order)
    # اگر بین این دو، سیستم crash کند چی؟
    
# راه‌حل: Outbox
async def create_order_GOOD(data):
    async with db.transaction():
        order = await db.create_order(data)
        # event در همان transaction
        await db.insert_outbox({
            "topic": "order.placed",
            "payload": order,
            "created_at": now(),
        })
    # یک background worker جدا، outbox را می‌خواند و publish می‌کند

# Outbox Worker
async def outbox_worker():
    while True:
        events = await db.fetch_unpublished_events(limit=100)
        for event in events:
            try:
                await rabbitmq.publish(event.topic, event.payload)
                await db.mark_published(event.id)
            except Exception:
                # next iteration retry
                pass
        await asyncio.sleep(1)
    

۲. Saga Pattern (preview – فصل بعد)

برای تراکنش‌های توزیع‌شده. در فصل ۷ به تفصیل بررسی می‌شود.

۳. Idempotent Consumer

پیام ممکن است چند بار delivered شود. consumer باید بتواند duplicate را تشخیص دهد.


async def handle_message(message):
    message_id = message["id"]
    
    # بررسی تکرار با Redis
    is_new = await redis.set(
        f"processed:{message_id}",
        "1",
        nx=True,  # فقط اگر وجود ندارد
        ex=86400  # ۲۴ ساعت TTL
    )
    
    if not is_new:
        logger.info(f"Duplicate {message_id}, skipping")
        return
    
    # پردازش
    await process(message)
    

۴. Inbox Pattern

برخی موارد می‌خواهیم event را در DB ذخیره کنیم و بعد پردازش کنیم (برای audit و idempotency).


async def consume_message(message):
    async with db.transaction():
        # ذخیره در inbox (UNIQUE message_id)
        try:
            await db.insert_inbox({
                "id": message["id"],
                "type": message["event_type"],
                "payload": message,
                "received_at": now(),
                "status": "pending"
            })
        except UniqueViolation:
            return  # duplicate
        
        # پردازش
        await process(message)
        
        # تأیید
        await db.update_inbox(message["id"], status="processed")
    

۶.۹ بهترین تجربیات

  1. Idempotent Consumer. فرض کنید پیام ممکن است چند بار delivered شود.
  2. Outbox Pattern. dual-write را اجتناب کنید.
  3. Schema Registry. برای evolution امن.
  4. Dead Letter Queue. پیام‌های شکست‌خورده را گم نکنید.
  5. Retry با Exponential Backoff. از hammering جلوگیری کنید.
  6. Message Versioning. فیلد version در هر event.
  7. Tracing. trace_id در headers قرار دهید.
  8. Monitoring. queue size، consumer lag، processing time.
  9. Backpressure. prefetch count و QoS تنظیم کنید.
  10. Persistent messages. برای پیام‌های مهم.
  11. Acknowledgment. manual ack بعد از موفقیت.
  12. Compression. snappy یا gzip برای کاهش bandwidth.
  13. Partition strategy. در Kafka، key مناسب برای ordering.
  14. Connection pooling. connection ها expensive هستند.
  15. Documentation. هر event مستند شود (purpose، schema، consumer ها).

۶.۱۰ خلاصه فصل

آنچه آموختیم:
  • الگوهای پیام‌رسانی: Point-to-Point، Pub/Sub، Request/Reply، Routing، Priority
  • RabbitMQ پیشرفته: production-grade publisher/consumer با DLQ و retry
  • Kafka پیشرفته: Producer با idempotence، Consumer Group با manual commit
  • Redis Pub/Sub و Streams برای real-time messaging
  • RabbitMQ vs Kafka: قانون «task = RabbitMQ، event = Kafka»
  • Schema Registry و evolution امن
  • الگوهای کلیدی: Outbox، Inbox، Idempotent Consumer
در فصل بعد: الگوی Saga — مدیریت تراکنش‌های توزیع‌شده در میکروسرویس‌ها. Choreography vs Orchestration، Compensation، و پیاده‌سازی عملی.

نمایش سایت

رنگ سایت
حالت نمایش
اندازهٔ متن
خوانایی

این تنظیمات فقط روی مرورگر شما ذخیره می‌شود.