מיקרוסרוויסים גדלים, ובקשות HTTP סינכרוניות הופכות את המערכת לרשת. כשל אחד — וכל המחסנית קורסת. ארכיטקטורה מונעת אירועים (EDA) היא הדרך היחידה לנתק את הקישוריות הזו. יישמנו EDA עבור פרויקטים עם עומסים של עד 10,000 אירועים בשנייה ואנחנו יודעים כיצד להימנע מטעויות נפוצות.
לקוח עם חנות מקוונת המעבדת 50,000 הזמנות ביום נתקל בפסקי זמן (timeouts) בעומסי שיא, שחסמו את המלאי וההודעות. לאחר יישום EDA עם Apache Kafka ותבנית ה-Outbox, זמן התגובה ירד ב-40%, התפוקה שולשה, ועלויות התמיכה ירדו ב-30% (מעל 5,000 דולר לחודש). זמן האחזור הממוצע לעיבוד אירוע הוא 10ms.
למה EDA עדיפה על אינטראקציה סינכרונית?
קריאות HTTP סינכרוניות הן כמו שיחות טלפון: אתה צריך תשובה מיידית. EDA עובדת כמו דואר: שלח מכתב ושכח. השוואה:
| פרמטר | סינכרוני | אסינכרוני (EDA) |
|---|---|---|
| תלות בזמן | הלקוח ממתין | הלקוח לא חסום |
| סובלנות לתקלות | כשל שובר את השרשרת | התור מבודד כשלים |
| עומס | יחס ישיר ל-RPS | חיץ + לחץ אחורי (backpressure) |
| מורכבות פיתוח | קל יותר להבנה | ניפוי באגים קשה יותר (אידמפוטנטיות, ניטור) |
איך אנו מיישמים EDA: מחסנית טכנולוגית ומקרה בוחן
אנו משתמשים ב-Apache Kafka כברוקר הראשי בשל התפוקה הגבוהה והאחסון לטווח ארוך.
מבנה אירוע
interface DomainEvent<T = unknown> { id: string; // UUID — для идемпотентности type: string; // 'user.registered', 'order.placed' version: string; // '1.0' — для schema evolution source: string; // 'order-service' correlationId: string; // сквозной ID через все сервисы causationId?: string; // ID события, ставшего причиной occurredAt: string; // ISO 8601 data: T; } // Конкретное событие interface OrderPlacedEvent extends DomainEvent<{ orderId: string; customerId: string; items: Array<{ productId: string; quantity: number; price: number }>; total: number; shippingAddress: Address; }> { type: 'order.placed'; } Apache Kafka — ברוקר ראשי
import { Kafka, Partitioners } from 'kafkajs'; const kafka = new Kafka({ clientId: 'order-service', brokers: process.env.KAFKA_BROKERS.split(',') }); // Продюсер const producer = kafka.producer({ createPartitioner: Partitioners.LegacyPartitioner }); async function publishOrderPlaced(order: Order): Promise<void> { await producer.send({ topic: 'order.events', messages: [{ key: order.id, // партиционирование по ID заказа value: JSON.stringify({ id: uuidv4(), type: 'order.placed', version: '1.0', source: 'order-service', correlationId: context.correlationId, occurredAt: new Date().toISOString(), data: { orderId: order.id, customerId: order.customerId, items: order.items, total: order.total } } satisfies OrderPlacedEvent), headers: { 'content-type': 'application/json', 'schema-version': '1.0' } }] }); } // Консьюмер — Inventory Service const consumer = kafka.consumer({ groupId: 'inventory-service' }); await consumer.subscribe({ topics: ['order.events'], fromBeginning: false }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { const event = JSON.parse(message.value.toString()) as DomainEvent; // Идемпотентность: проверяем, не обрабатывали ли уже это событие const processed = await idempotencyRepo.exists(event.id); if (processed) return; try { if (event.type === 'order.placed') { await inventoryService.reserveStock(event.data.orderId, event.data.items); } await idempotencyRepo.mark(event.id); } catch (error) { // Публикуем в Dead Letter Topic для анализа await deadLetterProducer.send({ topic: 'order.events.dlq', messages: [{ value: message.value, headers: { 'failure-reason': error.message } }] }); } } }); תבנית Outbox — אספקה מובטחת
טעות נפוצה: שמירה ל-DB ואז פרסום ל-Kafka — סיכון לאובדן אירועים. הגישה הנכונה — Outbox טרנזקציונלי:
// В рамках одной транзакции БД async function createOrder(dto: CreateOrderDto): Promise<Order> { return db.transaction(async (trx) => { // 1. Сохраняем заказ const order = await trx('orders').insert({ ...orderData }).returning('*'); // 2. Сохраняем событие в outbox-таблицу (в той же транзакции!) await trx('outbox_events').insert({ id: uuidv4(), aggregate_id: order.id, event_type: 'order.placed', payload: JSON.stringify(orderPlacedEvent), status: 'pending', created_at: new Date() }); return order; }); } קורא Outbox נפרד (Outbox Poller) קורא אירועים ממתינים ומפרסם אותם ל-Kafka:
// Cron job или background worker async function processOutbox(): Promise<void> { const events = await db('outbox_events') .where({ status: 'pending' }) .orderBy('created_at') .limit(100) .forUpdate() .skipLocked(); for (const event of events) { try { await kafka.producer.send({ topic: getTopicForEventType(event.event_type), messages: [{ key: event.aggregate_id, value: event.payload }] }); await db('outbox_events') .where({ id: event.id }) .update({ status: 'published', published_at: new Date() }); } catch { await db('outbox_events') .where({ id: event.id }) .update({ retry_count: db.raw('retry_count + 1') }); } } } חלופה — Debezium CDC: קורא את ה-WAL של PostgreSQL ומפרסם שינויים ל-Kafka ללא קוד.
איך להבטיח אספקה ללא אובדן?
טכניקות מפתח:
- Outbox טרנזקציונלי
- אידמפוטנטיות של הצרכן (בדיקה לפי event.id)
- תור הודעות מתות (Dead Letter Queue) להודעות שנכשלו
- ניטור זמן אחזור באמצעות Prometheus + Grafana
הקמנו מערכת המטפלת ב-1000 הודעות בשנייה ללא אובדן במהלך כשל של צומת Kafka בודד.
תיאום יעיל: כוריאוגרפיה מול אורכסטרציה
| מאפיין | כוריאוגרפיה | אורכסטרציה |
|---|---|---|
| תיאום | שירותים מגיבים לאירועים | אורכסטרטור מרכזי |
| צימוד | נמוך | בינוני |
| נראות זרימה | קשה למעקב | מפורש בקוד |
| בדיקות | קשה יותר | קל יותר |
כוריאוגרפיה מתאימה לתרחישים עם צימוד נמוך; אורכסטרציה כאשר נדרש רצף קפדני. אנו משלבים: בתוך שירות — אורכסטרציה, בין שירותים — כוריאוגרפיה.
מה זה Event Sourcing ו-CQRS?
ב-Event Sourcing, כל השינויים נשמרים כאירועים, המתפרסמים גם למאגר האירועים וגם לברוקר. CQRS (Command Query Responsibility Segregation) מפריד בין פעולות כתיבה וקריאה, ולעיתים קרובות משתמש באותם אירועים לבניית תחזיות (projections). EDA ו-CQRS/ES עובדים היטב יחד, אך דורשים תכנון מוקפד.
דוגמה לתצורת Kafka בייצור
להגדרה אמינה בייצור, הגדר גורם שכפול (replication factor) של 3, שימור (retention) של 7 ימים, ודחיסה (compaction) עבור נושאים קריטיים. מינימום 3 ברוקרים, ונטר את הפיגור של הצרכנים באמצעות Burrow. לסביבות ענן — Confluent Cloud או AWS MSK.
מה כלול בעבודת יישום EDA?
בעת הזמנה, אנו מספקים:
- תיעוד ארכיטקטוני (דיאגרמות, בחירת ברוקר)
- פיתוח יצרנים וצרכנים עם תבנית Outbox
- הגדרת Kafka/RabbitMQ עם התמדה ושכפול
- יישום אידמפוטנטיות ותור הודעות מתות (DLQ)
- בדיקות עומס עד 1000 הודעות בשנייה ללא אובדן
- ניטור והתראות (זמן אחזור, תפוקה)
- הכשרת צוות
קבל ייעוץ ליישום EDA עבור הפרויקט שלך — נעריך את הארכיטקטורה שלך ונציע פתרון אופטימלי.
כמה זמן לוקח יישום EDA?
| שלב | משך |
|---|---|
| תרחיש בודד (יצרן + 2–3 צרכנים) | 1–2 שבועות |
| + תבנית Outbox, אידמפוטנטיות, DLQ | +שבוע |
| EDA מלא עבור 5–10 שירותים + ניטור | 4–8 שבועות |
בחירת הברוקר תלויה בעומס ובדרישות. Kafka — לתפוקה גבוהה ואחסון ארוך. RabbitMQ — לניתוב מורכב. Redis Streams — לזמן אחזור נמוך. Google Pub/Sub — לענן. ברוב הפרויקטים, אנו משתמשים ב-Kafka.
יישום EDA ב-4 שלבים
- ניתוח ועיצוב: זיהוי גבולות הקשר, הגדרת אירועים.
- הגדרת ברוקר: פריסת Kafka עם שכפול, הגדרת נושאים (topics).
- יישום יצרנים וצרכנים: כתיבת קוד עם Outbox ואידמפוטנטיות.
- ניטור וניפוי באגים: הגדרת מדדים, DLQ, התראות.
לצוות שלנו ניסיון של 10+ שנים במערכות מבוזרות, עם למעלה מ-20 פרויקטי EDA. הזמן יישום EDA במפתחות מלאה — קבל ארכיטקטורה אמינה וניתנת להרחבה. צור קשר כדי להעריך את הפרויקט שלך.







