~/icsd.ir — bash
SYSTEM_ONLINE

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"])
    

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

  1. Event ها immutable. هیچ‌گاه تغییر نکنند.
  2. Self-contained events. consumer query نزند.
  3. Schema versioning. از روز اول.
  4. Past tense naming. OrderPlaced نه PlaceOrder.
  5. Outbox pattern. برای dual-write problem.
  6. Idempotent consumers. پیام تکراری مشکل ایجاد نکند.
  7. Snapshot برای event sourcing. performance.
  8. Read your own writes. در UI لحاظ کنید.
  9. Monitoring lag. consumer lag مهم است.
  10. 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 پیچیده
در فصل بعد: پایداری و Resilience — Circuit Breaker، Retry، Timeout، Bulkhead، Fault Tolerance.

نمایش سایت

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

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