Event-Driven Architecture
معماری Event-Driven (EDA) سبکی است که در آن سرویسها از طریق انتشار و گوش دادن به event ها با هم ارتباط برقرار میکنند. این الگو پایه بسیاری از سیستمهای مدرن مقیاسپذیر است.
۸.۱ مقدمه
معماری Event-Driven (EDA) سبکی است که در آن سرویسها از طریق انتشار و گوش دادن به event ها با هم ارتباط برقرار میکنند. این الگو پایه بسیاری از سیستمهای مدرن مقیاسپذیر است.
هدف فصل: درک Event-Driven Architecture، Event Sourcing، CQRS، Eventual Consistency و پیادهسازی عملی.
۸.۲ Event-Driven Architecture چیست؟
در EDA، سیستمها به جای فراخوانی مستقیم یکدیگر، از طریق event ها ارتباط برقرار میکنند. event یک «حقیقت» در گذشته است که اتفاق افتاده.
اصول EDA
- Event: یک facts immutable در گذشته (مثل
OrderPlaced) - Event Producer: سرویسی که event را منتشر میکند
- Event Consumer: سرویسی که به event واکنش نشان میدهد
- Event Bus / Broker: واسط بین producer و consumer (Kafka، RabbitMQ)
- Loose Coupling: producer نمیداند چه کسی consume میکند
تفاوت Command با Event
| Command | Event |
|---|---|
| «این کار را انجام بده» | «این کار اتفاق افتاده» |
| زمان حال (PlaceOrder) | زمان گذشته (OrderPlaced) |
| یک گیرنده | چند گیرنده |
| میتواند رد شود | قابل رد نیست (واقعیت است) |
| Coupling بیشتر | Loose coupling |
۸.۳ طراحی Event ها
ساختار یک Event خوب
{
"event_id": "evt_01HX2K7B...",
"event_type": "OrderPlaced",
"event_version": "1.0",
"timestamp": "2026-05-01T10:30:00Z",
"source": "order-service",
"trace_id": "trace_abc123",
"data": {
"order_id": 12345,
"user_id": 678,
"items": [...],
"total": 1500000,
"currency": "IRR"
},
"metadata": {
"user_agent": "...",
"ip_address": "..."
}
}
قواعد طراحی
- Immutable: event هیچگاه تغییر نمیکند
- Self-contained: همه دادههای لازم را داشته باشد (consumer نباید مجبور به query شود)
- Unique ID: برای deduplication و tracking
- Timestamp: زمان وقوع
- Versioned: برای schema evolution
- Past tense name: OrderPlaced، نه PlaceOrder
دو نوع Event
📦 Event-Carried State Transfer
Event تمام state لازم را حمل میکند. consumer نیازی به query ندارد.
OrderPlaced { order_id, user_id, items[], total, ... }
✅ Decoupling کامل، Performance بهتر
❌ Event ها بزرگتر
🔔 Event Notification
Event فقط اطلاع میدهد. consumer برای جزئیات query میزند.
OrderPlaced { order_id }
✅ Event ها کوچک
❌ Coupling به sender، query overhead
۸.۴ Event Sourcing
به جای ذخیره وضعیت فعلی، تمام تغییرات (event ها) را ذخیره میکنیم. وضعیت فعلی با replay کردن event ها بازسازی میشود.
مقایسه با CRUD سنتی
CRUD سنتی:
account_balance: 5000 (وضعیت فعلی - تاریخچه گم میشود)
Event Sourcing:
AccountCreated { initial: 0 }
MoneyDeposited { amount: 10000 }
MoneyWithdrawn { amount: 3000 }
MoneyDeposited { amount: 2000 }
MoneyWithdrawn { amount: 4000 }
→ balance = sum(events) = 5000
مزایا:
- تاریخچه کامل
- Audit trail خودکار
- قابلیت "time travel"
- Rebuild از hر زمان
پیادهسازی ساده
from dataclasses import dataclass
from typing import List
import json
@dataclass
class Event:
event_id: str
aggregate_id: str
event_type: str
data: dict
timestamp: float
version: int
class EventStore:
"""ذخیره event ها در PostgreSQL"""
async def append(self, event: Event):
await db.execute("""
INSERT INTO events (event_id, aggregate_id, event_type,
data, timestamp, version)
VALUES ($1, $2, $3, $4, $5, $6)
""", event.event_id, event.aggregate_id, event.event_type,
json.dumps(event.data), event.timestamp, event.version)
async def get_events(self, aggregate_id: str) -> List[Event]:
rows = await db.fetch("""
SELECT * FROM events
WHERE aggregate_id = $1
ORDER BY version ASC
""", aggregate_id)
return [Event(**row) for row in rows]
# Account Aggregate
class Account:
def __init__(self, account_id: str):
self.account_id = account_id
self.balance = 0
self.version = 0
self.uncommitted_events: List[Event] = []
@classmethod
async def load(cls, account_id: str, event_store: EventStore):
"""بازسازی از event ها"""
account = cls(account_id)
events = await event_store.get_events(account_id)
for event in events:
account._apply(event)
return account
def _apply(self, event: Event):
"""اعمال یک event روی state"""
if event.event_type == "AccountCreated":
self.balance = event.data["initial"]
elif event.event_type == "MoneyDeposited":
self.balance += event.data["amount"]
elif event.event_type == "MoneyWithdrawn":
self.balance -= event.data["amount"]
self.version = event.version
def deposit(self, amount: float):
if amount <= 0:
raise ValueError("Amount must be positive")
event = Event(
event_id=str(uuid.uuid4()),
aggregate_id=self.account_id,
event_type="MoneyDeposited",
data={"amount": amount},
timestamp=time.time(),
version=self.version + 1,
)
self._apply(event)
self.uncommitted_events.append(event)
def withdraw(self, amount: float):
if amount > self.balance:
raise InsufficientFundsError()
event = Event(
event_id=str(uuid.uuid4()),
aggregate_id=self.account_id,
event_type="MoneyWithdrawn",
data={"amount": amount},
timestamp=time.time(),
version=self.version + 1,
)
self._apply(event)
self.uncommitted_events.append(event)
async def save(self, event_store: EventStore):
for event in self.uncommitted_events:
await event_store.append(event)
self.uncommitted_events.clear()
Snapshot برای Performance
اگر یک aggregate هزاران event داشته باشد، replay کند است. راهحل: Snapshot.
# هر ۱۰۰ event، یک snapshot ذخیره کن
if account.version % 100 == 0:
await snapshot_store.save(account)
# بازسازی: load snapshot + replay event های بعد از آن
async def load_with_snapshot(account_id):
snapshot = await snapshot_store.get_latest(account_id)
account = Account.from_snapshot(snapshot)
events = await event_store.get_events_after(account_id, snapshot.version)
for event in events:
account._apply(event)
return account
۸.۵ CQRS — Command Query Responsibility Segregation
جدا کردن مدل خواندن (Query) از مدل نوشتن (Command). هر کدام میتوانند DB، schema و حتی technology متفاوت داشته باشند.
┌─────────────────────────────────────┐
│ Application │
├─────────────────┬───────────────────┤
│ Command Side │ Query Side │
│ │ │
│ Commands │ Queries │
│ │ │ ▲ │
│ ▼ │ │ │
│ Aggregate │ Read Model │
│ (Write DB) │ (Read DB) │
│ PostgreSQL │ Elasticsearch │
│ │ │ │
│ └────► Events ──────► │
└─────────────────────────────────────┘
چرا CQRS؟
- Optimization متفاوت: write برای ACID، read برای performance
- Scalability مستقل: read معمولاً ۹۰٪ ترافیک است
- Schema بهینه: read DB میتواند denormalized باشد
- چند view: یک write، چندین read model
پیادهسازی
# Command Side - PostgreSQL
class CreateProductCommand:
name: str
price: float
category_id: int
@app.post("/products")
async def create_product(cmd: CreateProductCommand):
async with db.transaction():
product = await Product.create(**cmd.dict())
# Event publish
await publish_event("product.created", {
"id": product.id,
"name": product.name,
"price": product.price,
"category_id": product.category_id,
})
return product
# Query Side - Elasticsearch (denormalized)
@event_handler("product.created")
async def update_product_search_index(event):
# دریافت اطلاعات کامل با join های لازم
category = await get_category(event["category_id"])
# ایجاد document بهینه برای search
doc = {
"id": event["id"],
"name": event["name"],
"price": event["price"],
"category_id": event["category_id"],
"category_name": category["name"], # denormalized
"category_path": category["path"],
"search_text": f"{event['name']} {category['name']}",
}
await es.index(index="products", id=event["id"], body=doc)
# Query API
@app.get("/products/search")
async def search_products(q: str, category: int = None):
query = {
"query": {
"bool": {
"must": [{"match": {"search_text": q}}],
"filter": [{"term": {"category_id": category}}] if category else []
}
}
}
return await es.search(index="products", body=query)
چالش CQRS: Eventual Consistency
بین update در write DB و reflect در read DB، یک تأخیر کوچک وجود دارد. کاربر ممکن است داده قدیمی ببیند.
راهحلها:
- Read your own writes — query از write DB برای کاربر خودش
- Optimistic UI — UI فوراً نشان دهد، در background sync کن
- Version checking — version number را با هر read برگردان
۸.۶ Eventual Consistency
در سیستمهای توزیعشده، طبق CAP theorem، نمیتوان همزمان Consistency، Availability و Partition Tolerance را داشت. اکثر سیستمهای مدرن AP را انتخاب میکنند (Available + Partition Tolerant) و C را به «Eventual Consistency» تنزل میدهند.
سطوح Consistency
- Strong Consistency: همه گرهها همیشه same data میبینند (مثل DB transaction)
- Eventual Consistency: در نهایت همه گرهها sync میشوند، اما ممکن است موقتاً متفاوت باشند
- Causal Consistency: اگر A قبل از B بود، همه این ترتیب را میبینند
- Read-Your-Writes: کاربر خودش همیشه آخرین تغییرات خود را میبیند
طراحی برای Eventual Consistency
- UI طراحی کنید که delay را hide کند. مثلاً «در حال پردازش…»
- Idempotency. در صورت retry، duplicate نشود
- Versioning. برای detect کردن stale data
- Compensation. اگر eventual consistency به conflict منجر شد
- Conflict resolution. Last-Write-Wins یا CRDT
۸.۷ Event Store
یک پایگاه داده تخصصی برای ذخیره event ها. میتواند PostgreSQL ساده باشد یا ابزار تخصصی مثل EventStoreDB.
Schema در PostgreSQL
CREATE TABLE events (
event_id UUID PRIMARY KEY,
aggregate_id VARCHAR(255) NOT NULL,
aggregate_type VARCHAR(100) NOT NULL,
event_type VARCHAR(100) NOT NULL,
event_version INT NOT NULL,
payload JSONB NOT NULL,
metadata JSONB,
created_at TIMESTAMPTZ DEFAULT NOW(),
UNIQUE (aggregate_id, event_version) -- Optimistic locking
);
CREATE INDEX idx_aggregate ON events (aggregate_id, event_version);
CREATE INDEX idx_event_type ON events (event_type, created_at);
CREATE INDEX idx_created_at ON events (created_at);
-- جلوگیری از تغییر event ها
CREATE RULE no_update AS ON UPDATE TO events DO INSTEAD NOTHING;
CREATE RULE no_delete AS ON DELETE TO events DO INSTEAD NOTHING;
EventStoreDB
ابزار تخصصی برای event sourcing با features زیر:
- Optimistic concurrency
- Built-in projections
- Subscribe to events
- Snapshot support
- Stream-based organization
۸.۸ الگوهای پیشرفته EDA
۱. Event Notification
سرویسها فقط اطلاع میدهند. consumer ها برای جزئیات query میزنند.
۲. Event-Carried State Transfer
Event تمام state را حمل میکند. سرویسها local cache میسازند.
۳. Event Sourcing
(در بالا توضیح داده شد.)
۴. CQRS
(در بالا توضیح داده شد.)
۵. Event-Driven Saga
(فصل ۷ توضیح داده شد.)
۶. Materialized View
یک view از دادههای چند سرویس که با event ها بهروز میشود.
# Materialized View برای dashboard کاربر
@event_handler("order.placed")
async def update_user_stats(event):
await db.execute("""
UPDATE user_stats
SET total_orders = total_orders + 1,
total_spent = total_spent + $1,
last_order_at = $2
WHERE user_id = $3
""", event["total"], event["timestamp"], event["user_id"])
@event_handler("review.submitted")
async def update_user_stats_review(event):
await db.execute("""
UPDATE user_stats
SET total_reviews = total_reviews + 1
WHERE user_id = $1
""", event["user_id"])
۸.۹ بهترین تجربیات
- Event ها immutable. هیچگاه تغییر نکنند.
- Self-contained events. consumer query نزند.
- Schema versioning. از روز اول.
- Past tense naming. OrderPlaced نه PlaceOrder.
- Outbox pattern. برای dual-write problem.
- Idempotent consumers. پیام تکراری مشکل ایجاد نکند.
- Snapshot برای event sourcing. performance.
- Read your own writes. در UI لحاظ کنید.
- Monitoring lag. consumer lag مهم است.
- Audit trail. از مزایای رایگان EDA است.
۸.۱۰ خلاصه فصل
آنچه آموختیم:
- EDA: ارتباط از طریق event های immutable
- Event vs Command — معنا و کاربرد متفاوت
- Event Sourcing: ذخیره event ها به جای state
- CQRS: جدا کردن read و write برای optimization
- Eventual Consistency و CAP theorem
- Event Store در PostgreSQL یا EventStoreDB
- Materialized Views برای queries پیچیده