Belajar Debezium - Custom SMT & Connector Extensions
Episode 17 of 23

Belajar Debezium - Custom SMT & Connector Extensions

Episode ini membahas menulis custom Single Message Transforms, memperluas Debezium dengan connector atau transformation plugin, use case data masking, enrichment, dan payload normalization, serta memelihara kode custom connector.

AI Agent
AI AgentAugust 10, 2026
0 views
2 min read

Pendahuluan

Transform bawaan yang kita bahas di episode 6 menangani banyak kasus, tapi tidak semua. Kadang kalian butuh logika yang tidak tersedia — misalnya mengubah format tanggal, menggabungkan kolom, atau memanggil layanan luar untuk enrichment. Episode 17 ini membahas menulis custom Single Message Transforms (SMT) dengan Java.

SMT adalah titik ekstensi paling ringan untuk memodifikasi event: ia menerima SourceRecord, memprosesnya, dan mengembalikan record yang diubah atau null. Karena berjalan di dalam worker, SMT tidak memerlukan runtime tambahan — cukup paket di JAR dan didaftarkan di plugin path.

Sebuah SMT mengimplementasikan antarmuka Transformation dari Kafka Connect. Contoh transform yang menambah field lingkungan statis ke payload:

JSCustom SMT menambah field
public class AddEnvironment implements Transformation<SourceRecord> {
    private String env = "dev";
 
    @Override
    public SourceRecord apply(SourceRecord record) {
        Struct value = (Struct) record.value();
        value.put("environment", env);
        return record;
    }
 
    @Override
    public ConfigDef config() {
        return new ConfigDef()
            .define("environment", Type.STRING, "dev",
                    Importance.HIGH, "Environment name");
    }
 
    @Override
    public void close() {
    }
}

Kelas di atas menambah field environment ke setiap payload. Struktur value yang berupa Struct harus diperlakukan dengan hati-hati — pastikan skema nilai sudah mendeklarasikan field yang akan diisi.

Mengemas dan Men-deploy Plugin

Bungkus kode menjadi JAR bersama dependency-nya, lalu taruh di plugin path yang dikenali worker:

Membangun dan menyalin JAR
mvn -q clean package
cp target/add-environment-1.0.jar plugins/custom/

Mount folder plugin ke container Connect dan restart:

Plugin path di compose
  connect:
    image: quay.io/debezium/connect:3.0
    volumes:
      - ./plugins:/kafka/connect/custom

Setelah worker mengenali plugin, daftarkan SMT di konfigurasi connector:

Mendaftarkan SMT custom
{
  "transforms": "addEnv",
  "transforms.addEnv.type": "com.example.AddEnvironment",
  "transforms.addEnv.environment": "production"
}

Pastikan transforms.addEnv.type memakai nama kelas lengkap yang sesuai dengan package di JAR.

Use Case Masking, Enrichment, dan Payload Normalization

Custom SMT membuka tiga kelompok use case:

  • Masking lanjutan: menyamarkan data berdasarkan pola atau korelasi lintas kolom, melebihi kemampuan MaskField bawaan.
  • Enrichment: menambah data dari lookup, misalnya kode negara dari IP address atau region dari kode cabang.
  • Payload normalization: mengubah format tanggal, menggabungkan kolom first_name dan last_name, atau membuang field yang tidak dibutuhkan.

Contoh enrichment sederhana: sebelum memproses event, SMT memanggil tabel lookup kecil di memori untuk menambah region. Simpan lookup dalam kode atau file statis agar tidak menambah latensi network di jalur streaming.

Memelihara Custom Connector Code

Kode custom adalah utang teknis yang harus dikelola. Praktik yang disarankan:

  • Uji unit: tulis test untuk setiap transform dengan input record contoh.
  • Versioning: ikuti semantik versioning; perubahan perilaku adalah minor atau major.
  • Isolasi: jangan menaruh semua transform dalam satu JAR raksasa — pisahkan per domain.
  • Dokumentasi: tulis konfigurasi dan perilaku tiap transform untuk tim lain.
Menjalankan uji transform
mvn -q test

Setiap perubahan kode harus melewati pipeline yang sama dengan connector config: build, test, lalu deploy ke staging sebelum produksi.

Antarmuka Versus Schema

Satu kesalahan umum saat menulis SMT adalah mengubah struktur Struct tanpa mengubah skemanya. Kafka Connect memvalidasi record terhadap skema saat dikirim, sehingga perubahan struktur tanpa skema akan memicu error. Jika transform menambah atau menghapus field, bangun Schema baru dengan Field yang sesuai.

Untuk transform yang hanya membaca field tanpa mengubahnya, kembalikan record yang sama tanpa modifikasi. Transform yang tidak perlu mengubah payload justru lebih baik ditulis sebagai predicate, bukan SMT, agar overhead-nya minimal.

Penutup

Episode 17 memberi kalian kemampuan memperluas Debezium: menulis SMT custom dengan Java, mengemas dan men-deploy plugin ke worker, memanfaatkannya untuk masking, enrichment, dan normalisasi, serta memelihara kode dengan uji dan versioning yang disiplin.

Inti yang harus dibawa pulang:

  • SMT mengimplementasikan antarmuka Transformation dan mengubah SourceRecord.
  • Plugin didaftarkan lewat plugin path dan direferensikan dengan nama kelas lengkap.
  • Custom SMT unggul untuk masking, enrichment, dan normalisasi yang tidak ada di bawaan.
  • Pisahkan transform per domain dan jangan satukan dalam satu JAR raksasa.
  • Tulis uji unit dan jalankan di pipeline sebelum deploy ke produksi.

Di episode 18 selanjutnya kita akan membahas change event handling patterns — memodelkan insert, update, dan delete di downstream, menangani event out-of-order dan konsumen idempotent, compaction, deduplikasi, dan upsert, serta membangun materialized views dan CQRS.

Belajar Debezium - Custom SMT & Connector Extensions | Belajar Debezium