CDC ile Real-Time Data Streaming

Veritabanındaki değişiklikleri yakalayıp (CDC) Kafka gibi sistemlere aktarmak için Debezium kullanımı, veri senkronizasyonu ve event-driven mimari entegrasyonu anlatılır.

CDC ile Real-Time Data Streaming

Change Data Capture (CDC) ve Debezium ile Gerçek Zamanlı Veri Akışı

Mikroservis mimarileri ve event-driven sistemlerin yükselişiyle birlikte, veritabanlarındaki her değişikliği (ekleme, güncelleme, silme) anında yakalayıp diğer sistemlere iletme ihtiyacı doğdu. Change Data Capture (CDC) tam olarak bunu yapar: Veritabanı işlem günlüklerini (transaction logs, WAL, binlog) okuyarak veri değişikliklerini gerçek zamanlı olarak tespit eder ve harici sistemlere iletir.

Debezium, bu işi yapmak için geliştirilmiş, açık kaynaklı bir CDC platformudur. Kafka Connect üzerinde çalışan Debezium, veritabanındaki değişiklikleri yakalar ve bunları Kafka topic'lerine akış olarak (stream) gönderir. Bu sayede, veritabanı ile diğer servisler arasında gerçek zamanlı, güvenilir ve ölçeklenebilir bir veri köprüsü kurulur.


1. Debezium Nasıl Çalışır?

Debezium, her veritabanının kendi değişiklik yakalama mekanizmasını kullanır:

  • PostgreSQL: WAL (Write-Ahead Log) üzerinden logical decoding ile değişiklikleri okur.

  • MySQL: Binary log'u (binlog) ROW formatında okuyarak her satırdaki değişikliği yakalar.

  • SQL Server: Transaction log'u okuyarak veya SQL Server'ın native CDC özelliğini kullanarak değişiklikleri tespit eder.

  • MongoDB: Oplog (operation log) üzerinden değişiklikleri okur.

Debezium, bir Kafka Connect bağdaştırıcısı (connector) olarak çalışır. Yani Debezium connector'ları, Kafka Connect framework'ü içinde çalıştırılır. Veritabanına bağlanır, log'ları okur ve her değişiklik için bir CDC event (genellikle JSON formatında) oluşturup bunu belirlenen Kafka topic'ine yazar.


2. Debezium ile Oluşturulan CDC Event'in Yapısı

Debezium'un Kafka'ya gönderdiği her mesaj (event), değişikliğin detaylarını içeren zengin bir JSON nesnesidir. Tipik bir event aşağıdaki gibidir:

json

{
  "schema": { ... },
  "payload": {
    "before": null,                    // Değişiklik öncesi veri (DELETE veya UPDATE için)
    "after": {                         // Değişiklik sonrası veri (INSERT veya UPDATE için)
      "id": 1001,
      "customer_name": "Ahmet Yılmaz",
      "total_amount": 1250.50,
      "status": "SHIPPED"
    },
    "source": {                        // Değişikliğin kaynağı hakkında metadata
      "version": "2.5.0.Final",
      "connector": "postgresql",
      "name": "inventory-connector",
      "ts_ms": 1700000000000,
      "snapshot": "false",
      "db": "ecommerce_db",
      "schema": "public",
      "table": "orders",
      "txId": 12345,
      "lsn": 987654321
    },
    "op": "u",                         // İşlem türü: c=create, u=update, d=delete, r=read
    "ts_ms": 1700000001000            // Event'in oluşturulduğu zaman
  }
}

Bu event, bir siparişin (orders tablosu) status alanının güncellendiğini (op: "u") ve güncel veriyi (after) içermektedir. source kısmı ise hangi veritabanı, tablo ve transaction'dan geldiği gibi önemli metadata bilgilerini taşır.


3. Debezium + Kafka Mimarisi

Tipik bir CDC pipeline'ı şu bileşenlerden oluşur:

  1. Kaynak Veritabanı (Source Database): Değişikliklerin gerçekleştiği işlemsel veritabanı (PostgreSQL, MySQL, SQL Server, vb.).

  2. Debezium Connector: Kafka Connect üzerinde çalışan, veritabanı log'larını okuyan ve değişiklikleri Kafka'ya yazan bağdaştırıcı.

  3. Apache Kafka: CDC event'lerinin toplandığı, dayanıklı (durable) ve ölçeklenebilir mesaj kuyruğu. Her tablo için ayrı bir topic oluşturulabilir.

  4. Schema Registry (Opsiyonel): Event'lerin şemasını (Avro, Protobuf, JSON Schema) yönetmek ve versiyonlamak için kullanılır.

  5. Tüketici (Consumer) Uygulamalar: Kafka topic'lerinden event'leri okuyan ve işleyen servisler (mikroservisler, veri ambarı, cache, arama motoru, vb.).

Bu mimari, kaynak veritabanına herhangi bir ek yük bindirmeden, değişiklikleri gerçek zamanlı olarak diğer sistemlere iletir.


4. Kullanım Senaryoları (Ne Zaman Kullanılır?)

CDC ve Debezium, birçok farklı senaryoda kullanılabilir:

  • Mikroservisler Arası Veri Senkronizasyonu: Bir servisin veritabanındaki değişiklikleri, diğer servislerin kendi veritabanlarına yansıtması. Örneğin, Sipariş Servisi'nde oluşan yeni sipariş, Stok Servisi'ne ve Kargo Servisi'ne anında bildirilir.

  • Event-Driven Mimari (EDA): Veritabanı değişikliklerini domain event'lere dönüştürerek, servislerin olaylara (event) tepki vermesini sağlamak.

  • Cache Invalidation (Önbellek Geçersizleştirme): Veritabanında bir güncelleme olduğunda, ilgili önbellek (cache) kaydını otomatik olarak geçersiz kılmak veya güncellemek.

  • Real-Time Analytics (Gerçek Zamanlı Analitik): Veritabanındaki değişiklikleri anında bir veri ambarına veya analitik sistemine (ör. Spark, Snowflake) akıtarak anlık raporlar oluşturmak.

  • Veri Replikasyonu ve Veri Göçü (Migration): Bir veritabanından diğerine kesintisiz veri kopyalamak veya farklı veritabanları arasında geçiş yapmak.

  • Audit Log / Değişiklik Geçmişi: Veritabanındaki tüm değişiklikleri (kim, ne zaman, ne değiştirdi) otomatik olarak loglamak.

  • Outbox Pattern: Servislerin veritabanına yazdığı event'leri (outbox tablosu) CDC ile yakalayıp Kafka'ya göndererek güvenilir event yayını sağlamak.


5. .NET ile Debezium Entegrasyonu

.NET uygulamaları, Debezium'un ürettiği Kafka event'lerini tüketebilir (consume) ve işleyebilir. Ayrıca, Debezium'un Outbox Pattern ile entegrasyonu, .NET mikroservislerinde veri tutarlılığını sağlamak için yaygın bir yaklaşımdır.

A. .NET'te Kafka Consumer ile CDC Event'lerini Okumak

Confluent.Kafka veya MassTransit gibi kütüphaneler kullanılarak Kafka topic'lerinden CDC event'leri okunabilir.

csharp

using Confluent.Kafka;
using System.Text.Json;

var config = new ConsumerConfig
{
    BootstrapServers = "localhost:9092",
    GroupId = "cdc-consumer-group",
    AutoOffsetReset = AutoOffsetReset.Earliest
};

using var consumer = new ConsumerBuilder<string, string>(config).Build();
consumer.Subscribe("dbserver1.public.orders");

while (true)
{
    var consumeResult = consumer.Consume();
    var cdcEvent = JsonDocument.Parse(consumeResult.Message.Value);

    var operation = cdcEvent.RootElement.GetProperty("payload").GetProperty("op").GetString();
    var after = cdcEvent.RootElement.GetProperty("payload").GetProperty("after");

    switch (operation)
    {
        case "c": // Create
            // Yeni kayıt eklendi
            var newOrder = JsonSerializer.Deserialize<Order>(after);
            await HandleOrderCreatedAsync(newOrder);
            break;
        case "u": // Update
            // Kayıt güncellendi
            var updatedOrder = JsonSerializer.Deserialize<Order>(after);
            await HandleOrderUpdatedAsync(updatedOrder);
            break;
        case "d": // Delete
            // Kayıt silindi
            var before = cdcEvent.RootElement.GetProperty("payload").GetProperty("before");
            var deletedOrder = JsonSerializer.Deserialize<Order>(before);
            await HandleOrderDeletedAsync(deletedOrder);
            break;
    }
}

B. Debezium Server ile .NET (Kafka'sız CDC Tüketimi)

Debezium, Kafka ihtiyacını ortadan kaldıran Debezium Server adında hafif bir mod da sunar. Debezium Server, CDC event'lerini doğrudan HTTP, NATS, Kinesis veya diğer hedeflere gönderebilir. Bu, özellikle küçük projeler veya Kafka kurmak istemeyen ekipler için idealdir.

yaml

# Debezium Server yapılandırması (application.properties)
debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector
debezium.source.offset.storage.file.filename=/data/offsets.dat
debezium.source.database.hostname=localhost
debezium.source.database.port=5432
debezium.source.database.user=postgres
debezium.source.database.password=postgres
debezium.source.database.dbname=ecommerce_db
debezium.source.table.include.list=public.orders

# Event'leri HTTP endpoint'e gönder
debezium.sink.type=http
debezium.sink.http.url=http://localhost:5000/api/cdc/events

.NET'te bir HTTP endpoint'i ile bu event'leri yakalayabilirsiniz.

C. Outbox Pattern ile Debezium Entegrasyonu ( .NET)

Outbox Pattern, veritabanı işlemi ile event yayınını atomik hale getirir. Servis, veritabanına hem iş verisini hem de bir outbox_event tablosuna event'i yazar. Debezium, outbox_event tablosundaki değişiklikleri yakalar ve Kafka'ya iletir.

csharp

// OrderService - Sipariş oluşturma metodu
public async Task CreateOrderAsync(Order order)
{
    using var transaction = await _context.Database.BeginTransactionAsync();
    try
    {
        // 1. Siparişi kaydet
        _context.Orders.Add(order);
        await _context.SaveChangesAsync();

        // 2. Outbox event'ini kaydet (aynı transaction içinde)
        var outboxEvent = new OutboxEvent
        {
            Id = Guid.NewGuid(),
            EventType = "OrderCreated",
            Payload = JsonSerializer.Serialize(order),
            CreatedAt = DateTime.UtcNow
        };
        _context.OutboxEvents.Add(outboxEvent);
        await _context.SaveChangesAsync();

        await transaction.CommitAsync();
    }
    catch
    {
        await transaction.RollbackAsync();
        throw;
    }
}

Bu yaklaşım, event'lerin kaybolmamasını ve tam olarak bir kez (exactly-once) işlenmesini garanti eder.


6. Debezium ile İlgili Sık Karşılaşılan Zorluklar

  • Veritabanı Log Yönetimi: Debezium, veritabanı log'larını (WAL, binlog) okur. Eğer Debezium uzun süre çalışmazsa, log'lar şişebilir ve disk alanı tükenebilir. Bu nedenle, Debezium'un sürekli çalışır durumda olması önemlidir.

  • Schema Evolution (Şema Değişiklikleri): Veritabanı şeması değiştiğinde (ör. yeni sütun eklendiğinde), Debezium'un bu değişikliği nasıl ele alacağını planlamak gerekir. Schema Registry kullanmak burada yardımcı olur.

  • İlk Yükleme (Snapshot): Debezium ilk çalıştırıldığında, mevcut tüm verileri bir snapshot olarak alır. Bu işlem büyük tablolarda uzun sürebilir ve veritabanına yük bindirebilir.

  • Event Sıralaması (Ordering): Kafka, aynı partition içinde event'lerin sırasını korur. Ancak, farklı partition'lardaki event'ler arasında sıralama garantisi yoktur. Bu nedenle, event'leri işlerken sıralama önemliyse, partition key'i doğru seçmek gerekir.

Sonuç:

Change Data Capture (CDC) ve Debezium, veritabanı ile event-driven dünya arasında güçlü ve güvenilir bir köprü kurar. Debezium, veritabanı değişikliklerini gerçek zamanlı olarak Kafka'ya akıtarak, mikroservisler arası senkronizasyon, gerçek zamanlı analitik, önbellek güncelleme ve outbox pattern gibi birçok senaryoyu mümkün kılar.

.NET ekosisteminde, Confluent.Kafka ile CDC event'lerini tüketebilir, Debezium Server ile Kafka'sız CDC çözümleri kurabilir veya Outbox Pattern ile Debezium'u entegre ederek güvenilir event yayını sağlayabilirsiniz. Debezium, veritabanı değişikliklerini gerçek zamanlı olarak yakalamanın en olgun ve yaygın kullanılan açık kaynak çözümüdür.

Tüm yazılar