پیامرسانی 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")
۶.۹ بهترین تجربیات
- Idempotent Consumer. فرض کنید پیام ممکن است چند بار delivered شود.
- Outbox Pattern. dual-write را اجتناب کنید.
- Schema Registry. برای evolution امن.
- Dead Letter Queue. پیامهای شکستخورده را گم نکنید.
- Retry با Exponential Backoff. از hammering جلوگیری کنید.
- Message Versioning. فیلد
versionدر هر event. - Tracing. trace_id در headers قرار دهید.
- Monitoring. queue size، consumer lag، processing time.
- Backpressure. prefetch count و QoS تنظیم کنید.
- Persistent messages. برای پیامهای مهم.
- Acknowledgment. manual ack بعد از موفقیت.
- Compression. snappy یا gzip برای کاهش bandwidth.
- Partition strategy. در Kafka، key مناسب برای ordering.
- Connection pooling. connection ها expensive هستند.
- 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