Data Engineering

پیاده‌سازی پایپ‌لاین‌های CDC بلادرنگ با Debezium و Kafka برای همگام‌سازی داده با تأخیر کم

در منظره داده‌های مدرن، پردازش دسته‌ای اغلب برای پاسخگویی به نیازهای تحلیل بلادرنگ، تشخیص تقلب و داشبوردهای زودرس کند است. راه حل در Capture داده‌های تغییر یافته (CDC) نهفته است؛ روشی که تغییرات سطح سطر را در پایگاه داده ردیابی کرده و آن‌ها را به سیستم‌های پایین‌دست منتقل می‌کند. با ترکیب Apache Kafka و Debezium، می‌توانید پایپ‌لاین‌های مقاوم و با تأخیر کم بسازید که اطمینان حاصل می‌کنند انبارهای داده و برنامه‌های شما همواره با سیستم‌های منبع همگام‌سازی شده‌اند.

چرا Debezium و Kafka؟

Debezium یک پلتفرم توزیع‌شده متن‌باز برای CDC است. این ابزار به عنوان یک سرویس اتصال‌دهنده برای Kafka Connect عمل می‌کند، لاگ‌های پایگاه داده (مانند Binlog برای MySQL، WAL برای PostgreSQL یا CDC برای SQL Server) را نظارت کرده و رویدادهای تغییر را به صورت موضوعات (Topics) Kafka منتشر می‌کند. Kafka به عنوان ستون فقرات بادوام و با ظرفیت بالا، این رویدادها را بافر می‌کند و به مصرف‌کنندگان متعدد اجازه می‌دهد داده‌ها را به صورت مستقل پردازش کنند بدون اینکه بر پایگاه داده منبع تأثیر بگذارند.

این معماری چندین مزیت کلیدی ارائه می‌دهد:

  • جداسازی (Decoupling): تولیدکنندگان (پایگاه‌های داده) و مصرف‌کنندگان (تحلیل‌ها، میکروسرویس‌ها) کاملاً از هم جدا شده‌اند.
  • قابلیت بازپخش (Replayability): از آنجا که Kafka داده‌ها را حفظ می‌کند، می‌توانید رویدادها را بازپخش کنید تا وضعیت را تعمیر کرده یا مصرف‌کنندگان جدید را اضافه کنید.
  • تأخیر کم: رویدادها در عرض میلی‌ثانیه نه ساعت‌ها یا روزها ضبط و تحویل داده می‌شوند.

راه‌اندازی اتصال‌دهنده Debezium

پیاده‌سازی CDC با پیکربندی یک اتصال‌دهنده Debezium آغاز می‌شود. اگرچه می‌توانید از REST API برای Kafka Connect استفاده کنید، اما بسیاری از استقرارها از فایل‌های پیکربندی یا ابزارهای مدیریت‌کننده مانند Strimzi استفاده می‌کنند. در زیر یک مثال عملی از پیکربندی JSON برای یک اتصال‌دهنده منبع MySQL آورده شده است.

اطمینان حاصل کنید که ثبت لاگ باینری (Binary Logging) MySQL فعال است و کاربر اختصاصی با امتیازات تکثیر (Replication) وجود دارد. پیکربندی زیر نحوه اتصال 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 منتشر می‌کند. این رویدادها از یک ساختار پوششی خاص پیروی می‌کنند که شامل متادیتا، تصویر قبل (Before-image) و تصویر بعد (After-image) از سطرها تغییر یافته است. برای نشان دادن همگام‌سازی با تأخیر کم، می‌توانید از یک مصرف‌کننده ساده Kafka در پایتون برای پردازش این رویدادها استفاده کنید.

در اینجا نحوه مصرف و تجزیه یک رویداد تغییر آمده است:

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"مشتری جدید: {payload['id']} - {payload['first_name']}")
    elif op == 'd':
        payload = value['before']
        print(f"مشتری حذف شده: {payload['id']}")

بهترین شیوه‌ها برای محیط تولید

برای اطمینان از پایداری در محیط تولید، موارد زیر را در نظر بگیرید:

  1. پایگاه داده طرح (Schema Registry): همیشه با Confluent Schema Registry یا مشابه آن یکپارچه شوید تا طرح‌های Avro/JSON را مدیریت کنید و از سازگاری با عقب‌گرد و ایمنی نوع اطمینان حاصل کنید.
  2. مدیریت خطا: صف‌های مرده (DLQ) را برای پیام‌هایی که پردازش آن‌ها شکست می‌خورد، پیکربندی کنید تا امکان بررسی دستی و تلاش مجدد فراهم شود.
  3. نظارت: معیارهای تأخیر (Lag) در Kafka و وضعیت اتصال‌دهنده را نظارت کنید. تأخیر بالا نشان می‌دهد که مصرف‌کنندگان پایین‌دست در به‌روز نگه داشتن خود مشکل دارند.

نتیجه‌گیری

پیاده‌سازی CDC بلادرنگ با Debezium و Kafka، پایگاه‌های داده ایستا را به جریان‌های داده پویا تبدیل می‌کند. این رویکرد نه تنها تأخیر را کاهش می‌دهد، بلکه انعطاف‌پذیری و مقیاس‌پذیری مورد نیاز برای مهندسی داده در سطح سازمانی را نیز فراهم می‌کند. با بهره‌گیری از این ابزارها، شما تصمیم‌گیری بلادرنگ را امکان‌پذیر کرده و یک منبع حقیقت واحد را در کل اکوسیستم داده خود حفظ می‌کنید.

Share: