Moderen veri dünyasında, toplu işleme genellikle gerçek zamanlı analiz, sahtekarlık tespiti ve canlı panoların gereksinimlerini karşılamak için yavaştır. Çözüm, veritabanındaki satır düzeyindeki değişiklikleri izleyen ve bunları aşağı akış sistemlere ileten Değişiklik Verisi Yakalama (CDC) yöntemindedir. Apache Kafka ve Debezium birleştirilerek, veri ambarlarınızın ve uygulamalarınızın kaynak sistemlerle her zaman senkronize kalmasını sağlayan sağlam, düşük gecikmeli hatlar oluşturabilirsiniz.
Neden Debezium ve Kafka?
Debezium, CDC için açık kaynaklı dağıtık bir platformdur. Kafka Connect için bir bağlayıcı hizmeti olarak çalışır; veritabanı günlüklerini (MySQL için Binlog, PostgreSQL için WAL veya SQL Server için CDC gibi) izler ve değişiklik olaylarını Kafka konularına yayar. Kafka, bu olayları arabellekleyen, kaynak veritabanını etkilemeden birden fazla tüketiciye bağımsız olarak veri işleme imkanı tanıyan dayanıklı ve yüksek verimli bir omurgadır.
Bu mimari birkaç temel avantaj sunar:
- Ayrıştırma (Decoupling): Üreticiler (veritabanları) ve tüketiciler (analitik, mikroservisler) tamamen birbirinden bağımsız hale getirilir.
- Tekrarlanabilirlik: Kafka veriyi sakladığından, durumu onarmak veya yeni tüketicileri sisteme dahil etmek için olayları yeniden oynatabilirsiniz.
- Düşük Gecikme: Olaylar saatler veya günler yerine milisaniyeler içinde yakalanır ve iletilir.
Debezium Bağlayıcısının Kurulumu
CDC uygulaması, bir Debezium bağlayıcısının yapılandırılmasıyla başlar. Kafka Connect REST API'sini kullanabilse de, birçok dağıtım yapılandırma dosyaları veya Strimzi gibi yönetim araçları kullanır. Aşağıda, bir MySQL kaynak bağlayıcısı için pratik bir JSON yapılandırma örneği verilmiştir.
MySQL ikili günlüklemenin (binary logging) etkin olduğundan ve replikasyon ayrıcalıklarına sahip özel bir kullanıcıdan olduğundan emin olun. Aşağıdaki yapılandırma, Debezium'un veritabanına nasıl bağlandığını ve tabloların Kafka konularına nasıl eşlendiğini tanımlar.
{
"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"
}
}
Değişiklik Olaylarını Tüketme ve İşleme
Bağlayıcı çalışmaya başladığında, Debezium olayları Kafka konularına yayar. Bu olaylar; meta veri, değişen satırların önceki görüntüsü (before-image) ve sonraki görüntüsünü (after-image) içeren belirli bir zarf yapısına uyar. Düşük gecikmeli senkronizasyonu göstermek için, bu olayları işlemek amacıyla Python'da basit bir Kafka tüketici kullanabilirsiniz.
Bir değişiklik olayını nasıl tüketip ayrıştırabileceğinize dair bir örnek:
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' create (oluştur), 'u' update (güncelle), 'd' delete (sil) için
if op in ['c', 'u']:
payload = value['after']
print(f"Yeni Müşteri: {payload['id']} - {payload['first_name']}")
elif op == 'd':
payload = value['before']
print(f"Silinecek Müşteri: {payload['id']}")
Üretim İçin En İyi Uygulamalar
Üretim ortamında kararlılığı sağlamak için aşağıdakileri göz önünde bulundurun:
- Şema Kaydı (Schema Registry): Geriye dönük uyumluluk ve tür güvenliğini sağlamak için her zaman Confluent Schema Registry veya benzeri bir araçla entegre olun ve Avro/JSON şemalarını yönetin.
- Hata Yönetimi: İşleme sırasında başarısız olan mesajlar için ölüm kuyrukları (DLQ) yapılandırın; bu sayede manuel inceleme ve yeniden deneme imkanı sağlanır.
- İzleme (Monitoring): Kafka'daki gecikme (lag) metriklerini ve bağlayıcı durumunu izleyin. Yüksek gecikme, aşağı akış tüketicilerinin geride kalmaya çalıştığını gösterir.
Sonuç
Debezium ve Kafka ile gerçek zamanlı CDC uygulamak, statik veritabanlarını dinamik veri akışlarına dönüştürür. Bu yaklaşım yalnızca gecikmeyi azaltmakla kalmaz, aynı zamanda kurumsal düzeyde veri mühendisliği için gereken dayanıklılığı ve ölçeklenebilirliği de sağlar. Bu araçlardan yararlanarak gerçek zamanlı karar vermeyi mümkün kılar ve tüm veri ekosisteminizde tek bir doğruluk kaynağını korursunuz.