Change Data Capture (CDC)
CDC adalah teknik untuk mendeteksi dan menangkap perubahan data (INSERT, UPDATE, DELETE) dari database sumber, lalu meneruskannya ke sistem lain secara real-time atau near real-time.
Mengapa CDC?
- Tanpa CDC — Full table scan setiap kali sync. Lambat untuk tabel besar, tidak tahu baris mana yang berubah.
- Dengan CDC — Hanya kirim perubahan (delta). Efisien, near real-time, tidak membebani database sumber.
Metode CDC
// 1. Timestamp-based CDC (paling sederhana)
// Butuh kolom updated_at di setiap tabel
SELECT * FROM orders
WHERE updated_at > '2024-06-15 10:00:00'; -- timestamp terakhir sync
// Kekurangan: tidak tangkap DELETE, butuh kolom timestamp
// 2. Trigger-based CDC
// Database trigger mencatat perubahan ke tabel log
CREATE TRIGGER orders_cdc_trigger
AFTER INSERT OR UPDATE OR DELETE ON orders
FOR EACH ROW EXECUTE FUNCTION log_change();
// Kekurangan: impact performance di database sumber
// 3. WAL-based CDC (recommended) ⭐
// Baca langsung dari Write-Ahead Log database
// Zero impact ke database sumber
// Tangkap semua operasi termasuk DELETE
Debezium
Debezium adalah platform CDC open-source yang membaca WAL dari database (PostgreSQL, MySQL, MongoDB) dan mengirim change events ke Kafka.
// Debezium change event (dari Kafka topic)
{
"before": { // State sebelum perubahan
"id": 1001,
"status": "pending",
"amount": 150000
},
"after": { // State setelah perubahan
"id": 1001,
"status": "paid",
"amount": 150000
},
"source": {
"connector": "postgresql",
"db": "myapp",
"table": "orders",
"ts_ms": 1718451234000
},
"op": "u" // c=create, u=update, d=delete, r=read(snapshot)
}
// Debezium connector config (JSON)
{
"name": "orders-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "db-host",
"database.port": "5432",
"database.user": "debezium",
"database.dbname": "myapp",
"table.include.list": "public.orders,public.users",
"topic.prefix": "myapp",
"plugin.name": "pgoutput"
}
}
CDC Pipeline Architecture
// Typical CDC setup
PostgreSQL ──WAL──→ Debezium ──→ Kafka ──→ Consumer ──→ Data Warehouse
│
├──→ Elasticsearch (search index)
├──→ Redis (cache invalidation)
└──→ Analytics service
// Use cases:
// 1. DB → Warehouse sync (tanpa full scan)
// 2. Cache invalidation (update cache saat DB berubah)
// 3. Search index sync (Elasticsearch)
// 4. Microservice event propagation
// 5. Audit log / compliance