الأنظمة الفورية والبنية التعاونية

الأنظمة الفورية وCRDTs ومعالجة التدفق

4 دقيقة للقراءة

ميزات الوقت الفعلي موجودة في كل مكان: المستندات التعاونية ولوحات المعلومات المباشرة والمحادثات والألعاب متعددة اللاعبين ومؤشرات الحضور. يغطي هذا الدرس البروتوكولات الأساسية وهياكل البيانات وأنماط المعالجة التي تشغّل هذه الأنظمة على نطاق واسع.

بروتوكولات النقل الفوري

توجد ثلاثة أساليب رئيسية لاتصال الخادم بالعميل. الاختيار الصحيح يعتمد على متطلبات زمن الاستجابة واتجاه الرسائل وقيود البنية التحتية.

يلجأ المرشحون إلى WebSocket تلقائيًا. وهو غالبًا الجواب الخاطئ: الاتصال المستمر ثنائي الاتجاه تكلفة دائمة في حالة الاتصال وتهيئة موازن الحمل وتعقيد النشر، ونصف الأنظمة التي يُطبَّق عليها لا تدفع البيانات إلا في اتجاه واحد.

أي وسيلة نقل تحتاجها هذه الميزة فعلًا؟

هل يحتاج العميل لإرسال رسائل متكررة، أم أنه يستقبل فقط؟

الميزةLong PollingSSE (أحداث مرسلة من الخادم)WebSocket
الاتجاهخادم إلى عميل (محاكاة)خادم إلى عميل فقطثنائي الاتجاه
البروتوكولHTTP/1.1 طلبات متكررةHTTP/1.1 مع text/event-streamTCP مُرقّى عبر ws:// أو wss://
عبء الاتصالمصافحة TCP جديدة لكل استطلاعاتصال مستمر واحداتصال مستمر واحد
البيانات الثنائيةلا (JSON/نص فقط)لا (UTF-8 نص فقط)نعم (إطارات ثنائية مدعومة)
إعادة الاتصال التلقائيتنفيذ يدويمدمج (واجهة EventSource)تنفيذ يدوي
توافق الوكلاء/CDNنعم (HTTP قياسي)نعم (HTTP قياسي)غالبًا يتطلب تكوين
مخاوف التوسعحجم طلبات عالٍاتصال واحد لكل عميلاتصال واحد لكل عميل
الأفضل لـالأنظمة القديمة، الإشعارات البسيطةالبث المباشر، لوحات المعلومات، مؤشرات الأسهمالمحادثات، التعاون، الألعاب

قاعدة عامة للمقابلة: استخدم SSE عندما يدفع الخادم البيانات باتجاه واحد (لوحات المعلومات، الإشعارات). استخدم WebSocket عندما يرسل العميل أيضًا بيانات متكررة (التحرير التعاوني، المحادثات). تجنب long polling في التصاميم الجديدة إلا عند وجود قيود بنية تحتية تمنع الاتصالات المستمرة.

إدارة اتصالات WebSocket على نطاق واسع

خادم واحد يمكنه الاحتفاظ بعشرات الآلاف من اتصالات WebSocket، لكن على نطاق واسع تواجه تحديات التوجيه وتجاوز الفشل وإدارة الحالة.

استراتيجيات توجيه الاتصالات

                    موازن الحمل
                    ┌───────────┐
  العميل A ────────>│           │──────> الخادم 1 (يحمل A, B)
  العميل B ────────>│  L7 LB    │──────> الخادم 2 (يحمل C, D)
  العميل C ────────>│           │──────> الخادم 3 (يحمل E, F)
  العميل D ────────>│           │
  العميل E ────────>│           │
  العميل F ────────>└───────────┘

الجلسات اللاصقة (Sticky sessions) تربط العميل بنفس الخادم باستخدام ملف تعريف ارتباط أو تجزئة معرّف العميل. بسيطة التنفيذ، لكنها تسبب مشاكل أثناء النشر المتدرج (الاتصالات تنقطع) وتوزيع حمل غير متساوٍ.

ترحيل الاتصال (Connection migration) يخزن حالة الاتصال خارجيًا (في Redis) حتى يتمكن أي خادم من استئناف الجلسة. عندما يتوقف خادم، يعيد العملاء الاتصال بأي خادم متاح يحمّل حالتهم من Redis. أكثر تعقيدًا لكن يزيل عيوب الجلسات اللاصقة.

نبضات القلب وإعادة الاتصال

import asyncio
import time
from dataclasses import dataclass, field

@dataclass
class ConnectionState:
    client_id: str
    last_heartbeat: float = field(default_factory=time.time)
    missed_heartbeats: int = 0

class HeartbeatManager:
    def __init__(self, interval_sec: float = 30.0, max_missed: int = 3):
        self.interval = interval_sec
        self.max_missed = max_missed
        self.connections: dict[str, ConnectionState] = {}

    def register(self, client_id: str) -> None:
        self.connections[client_id] = ConnectionState(client_id=client_id)

    def heartbeat_received(self, client_id: str) -> None:
        if client_id in self.connections:
            conn = self.connections[client_id]
            conn.last_heartbeat = time.time()
            conn.missed_heartbeats = 0

    def check_connections(self) -> list[str]:
        """إرجاع قائمة معرّفات العملاء التي يجب قطع اتصالها."""
        now = time.time()
        dead: list[str] = []
        for client_id, conn in self.connections.items():
            if now - conn.last_heartbeat > self.interval:
                conn.missed_heartbeats += 1
                if conn.missed_heartbeats >= self.max_missed:
                    dead.append(client_id)
        return dead

يجب على العملاء تنفيذ تراجع أسي مع اهتزاز (jitter) عند إعادة الاتصال لتجنب القطيع المتدافع عند إعادة تشغيل الخادم:

import random

def reconnect_delay(attempt: int, base_ms: int = 1000, max_ms: int = 30000) -> float:
    """تراجع أسي: 1 ثانية، 2، 4، 8... حد أقصى 30 ثانية، مع اهتزاز."""
    delay = min(base_ms * (2 ** attempt), max_ms)
    jitter = random.uniform(0, delay * 0.3)  # اهتزاز 0-30%
    return (delay + jitter) / 1000.0

بنية Pub/Sub

نمط Pub/Sub يفصل المنتجين عن المستهلكين، مما يتيح التوزيع لعدة مشتركين بدون اتصالات نقطة لنقطة.

Kafka مقابل Redis Streams

الجانبKafkaRedis Streams
نموذج التخزينسجل إلحاقي مقسّم على القرصتدفق في الذاكرة مع استمرارية اختيارية
الاحتفاظقابل للتكوين (أيام/أسابيع/للأبد)محدود بالذاكرة أو بطول أقصى
مجموعات المستهلكيننعم (تعيين مبني على التقسيمات)نعم (XREADGROUP مع إقرار)
ضمان الترتيبFIFO لكل تقسيمFIFO لكل تدفق
الإنتاجيةملايين الرسائل/ثانية (إدخال/إخراج مجمّع)مئات الآلاف/ثانية (خيط واحد)
زمن الاستجابةميللي ثوانٍ قليلة (التجميع يضيف تأخيرًا صغيرًا)أقل من ميللي ثانية
الأفضل لـمصادر الأحداث، خطوط تحليل البيانات، السجلات المتينةالإشعارات الفورية، pub/sub خفيف، أحداث طبقة التخزين المؤقت

أنماط التوزيع (Fan-Out)

التوزيع عند الكتابة (نموذج الدفع): عند نشر رسالة، يتم توصيلها فورًا لجميع المشتركين. يستخدمه Twitter للمستخدمين ذوي المتابعين القليلين. قراءات سريعة، كتابات مكلفة.

التوزيع عند القراءة (نموذج السحب): المشتركون يستعلمون عن الرسائل الجديدة عند الطلب. يستخدمه Twitter لحسابات المشاهير (ملايين المتابعين). كتابات رخيصة، قراءات أبطأ.

التوزيع عند الكتابة                التوزيع عند القراءة

المنتج                              المنتج
   │                                   │
   ▼                                   ▼
 نشر ──> صندوق المشترك A            سجل الموضوع
       ──> صندوق المشترك B              │
       ──> صندوق المشترك C           ┌──┼──┐
                                     ▼  ▼  ▼
  (N كتابة لكل رسالة)             A   B   C  (قراءة عند الطلب)
                                   (N قراءة لكل مستهلك)

رؤية للمقابلة: معظم الأنظمة تستخدم نهجًا هجينًا، والجزء المثير أن الحد الفاصل مقبض ضبط لا قرار معماري. ادفع للحسابات العادية، وتخطَّ التوزيع للحسابات كثيرة المتابعين، وادمج منشوراتها وقت القراءة. الكتابات العامة عن بنية الجدول الزمني في Twitter تضع الحد في حدود عشرة آلاف متابع، لكن الرقم نفسه أقل أهمية بكثير من كونه إعدادًا يعمل وقت التشغيل — يمكنك تحريكه تحت الحِمل، أو لكل منطقة، أو لكل فئة حسابات، دون إعادة نشر أي شيء.

سبب وجود النهج الهجين هو تضخّم الكتابة. منشور واحد من حساب له 50 مليون متابع يعني 50 مليون كتابة جدول زمني تحت التوزيع عند الكتابة الخالص؛ وعشرة منشورات يوميًا تعني نصف مليار. أما التوزيع عند القراءة فيكلّف عملية دمج واحدة لكل تحميل جدول زمني، وهو رقم أصغر بكثير، والأهم أنه محدود.

CRDTs للتحرير التعاوني

عندما يحرر عدة مستخدمين نفس المستند في وقت واحد، تحتاج استراتيجية لدمج التغييرات المتزامنة بدون تعارضات.

OT مقابل CRDTs

التحويل التشغيلي (OT) يحوّل العمليات ضد بعضها للحفاظ على الاتساق. Google Docs يستخدم OT مع خادم مركزي يرتب العمليات.

أنواع البيانات المنسوخة الخالية من التعارض (CRDTs) هي هياكل بيانات مصممة بحيث تتقارب التحديثات المتزامنة دائمًا بدون تنسيق.

تنبيه قبل أن تستشهد بـ Figma، فهذه تفصيلة سيلتقطها أي محاور قرأ المصدر. يُلخَّص مقال Figma الهندسي لعام 2019 كثيرًا بأن «Figma استبدلت OT بـ CRDTs»، وهذا التلخيص خاطئ في الجزء المهم منه. تقول مدونة Figma صراحة إنهم لا يستخدمون CRDTs حقيقية: «صُممت CRDTs لأنظمة لامركزية لا توجد فيها سلطة مركزية واحدة»، وFigma لديها واحدة — «خادمنا هو السلطة المركزية». استعاروا من CRDTs فكرة سجلات «آخر كاتب يفوز» وأسقطوا الآلية التي وُجدت فقط للتعايش مع غياب المنسّق، وهي بالضبط الآلية التي تجعل CRDTs مكلفة.1

النسخة الأمينة أنفع في المقابلة على أي حال: CRDTs تشتري لك تقاربًا بلا تنسيق، وإذا كان لديك خادم مركزي أصلًا فأنت تدفع ثمن ضمان لا تستخدمه.

OT مقابل CRDTs — والخيار الأوسط الذي تشحنه المنتجات فعليًا

Google Docs

التحويل التشغيلي (OT)

خادم مركزيمطلوب — هو من يرتب العمليات
بيانات وصفية لكل حرفلا شيء — عمليات فقط
التحرير دون اتصالمحدود
المزايا
  • عبء ذاكرة منخفض جدًا: المستند هو المستند، بلا هوية لكل حرف
  • الحذف يحذف فعلًا، فلا يتراكم في المستند طويل العمر ركام من الشواهد
  • عقود من الاستخدام الإنتاجي في Google Docs — حالاته الحدية معروفة
العيوب
  • دوال التحويل تتضاعف مع أنواع العمليات ويصعب إثبات صحتها إلى حد سيئ السمعة
  • خادم الترتيب إلزامي، فلا وجود لسيناريو ند لند
  • التحريرات دون اتصال تحتاج تحويلًا من الخادم عند العودة، ما يحد من مدة غياب العميل
Automerge، Yjs

CRDTs الحقيقية

خادم مركزيغير مطلوب
بيانات وصفية لكل حرفمعرّف فريد لكل عنصر
التحرير دون اتصالكامل — دمج عند العودة
المزايا
  • التقارب خاصية رياضية لا مجموعة اختبارات — التحريرات المتزامنة لا يمكن أن تتباعد
  • ند لند حقيقي: عميلان يتزامنان مباشرة بلا خادم
  • فترات انقطاع طويلة اعتباطيًا تُدمج بنظافة عند العودة
العيوب
  • هوية لكل عنصر تعني عبء ذاكرة حقيقيًا، والأحرف المحذوفة تبقى شواهد تحتاج تجميع نفايات
  • المستند ينمو مع تاريخ التحرير لا مع طول المحتوى
  • تدفع هذا كله مقابل تحرر من التنسيق لا تحتاجه إن كنت تملك خادمًا
ما بنته Figma

خادم مركزي بإلهام CRDT

خادم مركزينعم — هو من يحدد الترتيب
بيانات وصفية لكل حرفأُسقطت، ومعها الساعات المتجهة
التحرير دون اتصالمحدود
المزايا
  • دلالات «آخر كاتب يفوز» مستعارة من سجلات CRDT، والخادم هو الساعة
  • بلا ساعات متجهة وبلا تجميع نفايات للشواهد، فهو أخف وأسرع من CRDT حقيقية
  • أسهل في التعقل بكثير من مصفوفة تحويلات OT
العيوب
  • الخادم اعتماد صلب ونقطة فشل واحدة — لا ند لند
  • «آخر كاتب يفوز» يتخلص بصمت من تحرير الخاسر؛ هذا قرار منتج لا وجبة تقنية مجانية
  • تتنازل عن ضمان التقارب الشكلي الذي كان سبب وجود CRDTs أصلًا

CRDTs العدادات

G-Counter (عداد النمو فقط): كل عقدة تحتفظ بعدادها الخاص. القيمة هي مجموع جميع العقد. الدمج يأخذ الحد الأقصى لعداد كل عقدة.

from typing import Dict

class GCounter:
    """عداد CRDT للنمو فقط. يدعم الزيادة والدمج."""
    def __init__(self, node_id: str):
        self.node_id = node_id
        self.counts: Dict[str, int] = {node_id: 0}

    def increment(self, amount: int = 1) -> None:
        self.counts[self.node_id] = self.counts.get(self.node_id, 0) + amount

    def value(self) -> int:
        return sum(self.counts.values())

    def merge(self, other: "GCounter") -> None:
        for node_id, count in other.counts.items():
            self.counts[node_id] = max(self.counts.get(node_id, 0), count)

PN-Counter (عداد موجب-سالب): عدادي G-Counter — واحد للزيادات وواحد للنقصانات. القيمة = P.value() - N.value().

LWW-Register (سجل آخر كاتب يفوز)

يخزن قيمة واحدة مع طابع زمني. عند التعارض، الطابع الزمني الأعلى يفوز.

from dataclasses import dataclass
from typing import Any

@dataclass
class LWWRegister:
    """سجل آخر كاتب يفوز. الطابع الزمني الأعلى يفوز دائمًا."""
    value: Any = None
    timestamp: float = 0.0
    node_id: str = ""

    def set(self, value: Any, timestamp: float, node_id: str) -> None:
        if timestamp > self.timestamp or (
            timestamp == self.timestamp and node_id > self.node_id
        ):
            self.value = value
            self.timestamp = timestamp
            self.node_id = node_id

    def merge(self, other: "LWWRegister") -> None:
        if other.timestamp > self.timestamp or (
            other.timestamp == self.timestamp and other.node_id > self.node_id
        ):
            self.value = other.value
            self.timestamp = other.timestamp
            self.node_id = other.node_id

OR-Set (مجموعة الإضافة-الحذف المرصود)

تدعم الإضافة والحذف معًا. كل عنصر يُوسم بمعرّف فريد عند الإضافة. الحذف يزيل فقط الوسوم المحددة التي تمت ملاحظتها، لذا الإضافة المتزامنة لنفس العنصر تُحفظ.

CRDTs النصية: RGA و LSEQ

للتحرير التعاوني للنصوص، تحتاج CRDTs تنمذج تسلسلات الأحرف:

  • RGA (المصفوفة القابلة للنمو المنسوخة): كل حرف له معرّف فريد (node_id, sequence). الإدراج يضع حرفًا بعد موضع مرجعي. الحذف يُعلّم الأحرف كشواهد. الإدراجات المتزامنة في نفس الموضع تُرتب بشكل حتمي بمعرّف العقدة.

  • LSEQ: يُعيّن لكل حرف موضعًا من فضاء معرّفات كثيف. المواضع تُخصص بين الجيران الحاليين، متجنبة إعادة التوازن. أكثر كفاءة في المساحة من RGA للمستندات الكبيرة.

  مثال إدراج RGA (تحريرات متزامنة):

  البداية:     H - E - L - L - O

  المستخدم A يُدرج 'X' بعد 'E':   H - E - X - L - L - O
  المستخدم B يُدرج 'Y' بعد 'E':   H - E - Y - L - L - O

  بعد الدمج (A.id > B.id):        H - E - X - Y - L - L - O
  كلا الإدراجين محفوظان، ترتيب حتمي بمعرّف العقدة.

أنماط معالجة التدفق

التحليلات الفورية ومعالجة الأحداث تتطلب استراتيجيات تقسيم نوافذ لتجميع تدفقات البيانات غير المحدودة.

أنواع النوافذ

ثلاث طرق لتقطيع تدفق غير محدود إلى أجزاء منتهية

عرض ثابت، بلا تداخل. كل حدث ينتمي لنافذة واحدة بالضبط. |----5 دقائق----|----5 دقائق----|----5 دقائق----| | أحداث | أحداث | أحداث | الخيار الأرخص: حالة نافذة واحدة مفتوحة في كل لحظة. استخدمه للتجميع الدوري — مشاهدات الصفحة لكل شريحة خمس دقائق، وعدادات الفوترة بالساعة. الثمن هو العمى عند الحدود. الارتفاع المفاجئ الممتد بين 10:04 و10:06 ينقسم بين شريحتين وقد لا يبدو ارتفاعًا في أي منهما.
نوع النافذةحالة الاستخداممثال
متقلبةتجميع دوريمشاهدات الصفحة لكل شريحة 5 دقائق
منزلقةمقاييس متحركةمتوسط وقت الاستجابة خلال آخر 10 دقائق، يُحدّث كل دقيقة
الجلسةسلوك المستخدمإجمالي الإجراءات لكل جلسة مستخدم (تنتهي بعد 30 دقيقة خمول)

دلالات التسليم مرة واحدة بالضبط

في معالجة التدفق الموزعة، توجد ثلاثة ضمانات تسليم:

  • مرة واحدة على الأكثر: أطلق وانسَ. الرسائل قد تُفقد. الأسرع.
  • مرة واحدة على الأقل: أعد المحاولة عند الفشل. الرسائل قد تتكرر. يتطلب مستهلكين عديمي القوة.
  • مرة واحدة بالضبط: كل رسالة تُعالج مرة واحدة فقط. يُحقق عبر كتابات عديمة القوة + إزاحات معاملاتية (معاملات Kafka) أو إزالة التكرار بمعرّفات أحداث فريدة.

العلامات المائية (Watermarks)

العلامة المائية هي تأكيد طابع زمني: "جميع الأحداث ذات الطابع الزمني <= W قد وصلت." العلامات المائية تتيح للنظام معرفة متى يمكن إغلاق النافذة وإصدار النتائج، حتى مع الأحداث غير المرتبة.

وقت الحدث:    10:00  10:01  10:03  10:02  10:04  10:05
               │      │      │      │      │      │
العلامة المائية: ──────────────────────────────────>
               W=10:00       W=10:02       W=10:04

حدث متأخر (10:02 يصل بعد W=10:02) يمكن:
- إسقاطه (الأبسط)
- وضعه في مخرج جانبي لإعادة المعالجة
- تشغيل تحديث النافذة (تأخر مسموح)

كشف الحضور على نطاق واسع

عرض حالة "متصل" لملايين المستخدمين يتطلب تصميمًا دقيقًا لتجنب إثقال النظام.

البنية

الحيلة الجديرة بالتذكر أن أحدًا لا يكتب «غير متصل» في أي مكان. الانقطاع هو غياب كتابة، ما يعني أن العميل الذي ينهار أو تنفد بطاريته أو يدخل نفقًا يُعالَج بنفس الآلية التي تعالج من يسجّل خروجه بأدب.

الحضور عبر انتهاء المفتاح — الانقطاع لا يُكتب بل يُستنتج

يُضبط الـ TTL على ضعف فاصل النبضة تقريبًا، حتى لا تقلب حزمة ضائعة واحدة حالة المستخدم إلى غير متصل.

نبضةتجديد TTLبلا تجديدموجود؟أصبح غير متصلالعميلنبضة كل 30 ثانيةخدمة الحضورعديمة الحالة — أي مثيل يفي بالغ…TTL 60 ثانيةمفتاح RedisSET user:{id} EX 60انقطاع مُستنتَجانتهى المفتاحلم تصل نبضة في الوقت المحدداستعلام الحالةMGET user:101 user:102 ...قناة Pub/subادفع التغيّر للأصدقاء المتابعين

لعرض الحضور للأصدقاء/القنوات، تجنب الاستعلام لكل صديق. بدلاً من ذلك:

  1. Pub/sub لتغييرات الحالة: عندما تتغير حالة المستخدم (متصل/غير متصل)، انشر في قناة. المشتركون (الأصدقاء المتصلون حاليًا) يستلمون التحديث.
  2. التحميل الكسول: جلب الحضور فقط للمستخدمين المرئيين على الشاشة.
  3. الاستعلامات المجمّعة: MGET user:101 user:102 user:103 في استدعاء Redis واحد.

تطبيق المقابلة: تصميم محرر مستندات فوري

عند طلب تصميم Google Docs أو محرر تعاوني في مقابلة، هيكل إجابتك:

  1. النقل: WebSocket للتحرير الفوري ثنائي الاتجاه. رجوع إلى SSE + HTTP POST للشبكات المقيدة.
  2. حل التعارضات: CRDTs (RGA للنصوص) للدمج التلقائي بدون خادم ترتيب مركزي. اذكر OT كبديل (نهج Google Docs) والمقايضات.
  3. حالة المستند: كل حرف له معرّف CRDT فريد. العملاء يطبقون العمليات المحلية فورًا (متفائل) ويبثون للخادم. الخادم يدمج ويعيد البث.
  4. الحضور: مواضع المؤشرات وحالة المستخدم تُتتبع عبر نبضات القلب مع مفاتيح TTL في Redis.
  5. الاستمرارية: العمليات تُلحق بسجل (مصادر الأحداث). لقطات دورية للتحميل السريع للمستندات.
  6. التوسع: تقسيم المستندات عبر الخوادم. كل مستند يعيش على خادم واحد (توجيه لاصق بمعرّف المستند). العمليات عبر المستندات نادرة.

التالي: اختبار ومختبر الوحدة 03 — ابنِ واجهة خلفية لمستند تعاوني مع CRDTs وإدارة WebSocket وتتبع الحضور. :::

Footnotes

  1. إيفان والاس، "How Figma's multiplayer technology works"، مدونة Figma الهندسية، 16 أكتوبر 2019.

اختبار

اختبار الوحدة 3: الأنظمة الفورية والبنية التعاونية

خذ الاختبار
هل كان هذا الدرس مفيدًا؟

سجّل الدخول للتقييم