Data Engineering

تنفيذ خطوط أنابيب CDC في الوقت الفعلي باستخدام Debezium وKafka لمزامنة البيانات منخفضة التأخير

في مشهد البيانات الحديث، غالباً ما تكون معالجة الدفعات بطيئة جداً لتلبية متطلبات التحليل في الوقت الفعلي، وكشف الاحتيال، ولوحات المعلومات المباشرة. تكمن الحل في التقاط تغييرات البيانات (CDC)، وهي منهجية تتبع التغييرات على مستوى الصف في قاعدة البيانات ونشرها إلى الأنظمة التالية. من خلال الجمع بين Apache Kafka وDebezium، يمكنك بناء خطوط أنابيب قوية ومنخفضة التأخير تضمن أن مستودعات البيانات والتطبيقات الخاصة بك متزامنة دائماً مع أنظمة المصدر.

لماذا Debezium وKafka؟

Debezium هو منصة مفتوحة المصدر للتوزيع لتقنية CDC. يعمل كخدمة موصل لـ Kafka Connect، حيث يراقب سجلات قاعدة البيانات (مثل Binlog لـ MySQL، وWAL لـ PostgreSQL، أو CDC لـ SQL Server) ويصدر أحداث التغيير كعناوين Kafka. يعمل Kafka كعمود فقري دائم وعالي الإنتاجية يخزن هذه الأحداث، مما يسمح لمستهلكين متعددين بمعالجة البيانات بشكل مستقل دون التأثير على قاعدة البيانات المصدر.

تقدم هذه البنية عدة مزايا رئيسية:

  • فصل المكونات: يتم فصل المنتجين (قواعد البيانات) عن المستهلكين (التحليلات، الخدمات المصغرة) تماماً.
  • إمكانية إعادة التشغيل: نظراً لأن Kafka يحتفظ بالبيانات، يمكنك إعادة تشغيل الأحداث لإصلاح الحالة أو إضافة مستهلكين جدد.
  • تأخير منخفض: يتم التقاط الأحداث وتسليمها في أجزاء من الثانية، وليس بالساعات أو الأيام.

إعداد موصل Debezium

يبدأ تنفيذ CDC بتكوين موصل Debezium. بينما يمكنك استخدام واجهة برمجة التطبيقات REST لـ Kafka Connect، تستخدم العديد من عمليات النشر ملفات التكوين أو أدوات الإدارة مثل Strimzi. فيما يلي مثال عملي لتكوين JSON لموصل مصدر MySQL.

تأكد من تمكين تسجيل الثنائي لـ MySQL ولديك مستخدم مخصص يتمتع بصلاحيات النسخ. يحدد التكوين أدناه كيفية اتصال Debezium بقاعدة البيانات وكيفية تعيين الجداول إلى عناوين Kafka.

{
  "name": "mysql-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "mysql-host",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "dbz",
    "database.server.id": "184054",
    "database.server.name": "dbserver1",
    "database.include.list": "inventory",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.inventory",
    "include.schema.changes": "true",
    "snapshot.mode": "initial"
  }
}

استهلاك ومعالجة أحداث التغيير

بمجرد تشغيل الموصل، يصدر Debezium أحداثاً إلى عناوين Kafka. تتبع هذه الأحداث هيكل غلاف محدد يحتوي على البيانات الوصفية، والصورة السابقة، والصورة اللاحقة للصفوف التي تم تغييرها. لتوضيح المزامنة منخفضة التأخير، يمكنك استخدام مستهلك Kafka بسيط في لغة Python لمعالجة هذه الأحداث.

إليك كيفية استهلاك وتحليل حدث تغيير:

from kafka import KafkaConsumer
import json

consumer = KafkaConsumer(
    'dbserver1.inventory.customers',
    bootstrap_servers=['localhost:9092'],
    value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)

for msg in consumer:
    value = msg.value
    op = value['op']  # 'c' للإنشاء، 'u' للتحديث، 'd' للحذف
    
    if op in ['c', 'u']:
        payload = value['after']
        print(f"New Customer: {payload['id']} - {payload['first_name']}")
    elif op == 'd':
        payload = value['before']
        print(f"Deleted Customer: {payload['id']}")

أفضل الممارسات للإنتاج

لضمان الاستقرار في بيئة الإنتاج، ضع في الاعتبار ما يلي:

  1. سجل المخطط (Schema Registry): قم دائماً بالتكامل مع Confluent Schema Registry أو ما شابه لإدارة مخططات Avro/JSON، مما يضمن التوافق مع الإصدارات السابقة وأمان النوع.
  2. معالجة الأخطاء: قم بتكوين طوابير الرسائل الميتة (DLQs) للرسائل التي تفشل في المعالجة، مما يسمح بالفحص اليدوي وإعادة المحاولة.
  3. المراقبة: راقب مقاييس التأخير في Kafka وحالة الموصل. يشير التأخير العالي إلى أن المستهلكين في الطرف الآخر يواجهون صعوبة في مواكبة البيانات.

الخاتمة

يؤدي تنفيذ CDC في الوقت الفعلي باستخدام Debezium وKafka إلى تحويل قواعد البيانات الثابتة إلى تدفقات بيانات ديناميكية. لا يقلل هذا النهج من التأخير فحسب، بل يوفر أيضاً المرونة والقابلية للتوسع المطلوبة لهندسة البيانات على مستوى المؤسسات. من خلال الاستفادة من هذه الأدوات، تتيح اتخاذ القرارات في الوقت الفعلي وتحافظ على مصدر واحد للحقيقة عبر نظام البيانات بأكمله.

Share: