All Products
Search
Document Center

Data Transmission Service:Sinkronisasi data dari ApsaraDB for MongoDB ke Message Queue for Apache Kafka

Last Updated:Aug 28, 2026

Layanan Transmisi Data (DTS) memungkinkan Anda mengalirkan data perubahan dari ApsaraDB for MongoDB ke instans Message Queue for Apache Kafka. Langkah-langkah berikut menjelaskan cara membuat tugas sinkronisasi dengan menggunakan replica set MongoDB sebagai sumber dan instans Kafka sebagai tujuan.

Apa yang dicakup oleh tugas ini:

  • Arsitektur yang didukung: Replica set dan kluster sharded

  • Metode sinkronisasi inkremental: Oplog (direkomendasikan) dan Change Streams

  • Format data: Canal JSON yang dikirimkan ke topik Kafka

  • Cakupan sinkronisasi: Tingkat koleksi; sinkronisasi penuh dan inkremental

  • Penagihan: Sinkronisasi data penuh gratis; sinkronisasi data inkremental dikenai biaya

Prasyarat

Sebelum memulai, pastikan Anda telah memiliki:

Untuk versi database sumber dan tujuan yang didukung, lihat Ikhtisar solusi sinkronisasi.

Batasan

Tinjau batasan berikut sebelum membuat tugas sinkronisasi.

Batasan database sumber

  • Server sumber harus memiliki bandwidth keluar yang mencukupi. Bandwidth yang tidak mencukupi akan mengurangi kecepatan sinkronisasi.

  • Saat mengonfigurasi pemetaan nama untuk koleksi, satu tugas dapat menyinkronkan hingga 1.000 koleksi. Melebihi batas ini akan menyebabkan error permintaan. Untuk menyinkronkan lebih dari 1.000 koleksi, buat beberapa tugas atau sinkronkan seluruh database tanpa pemetaan nama.

  • Untuk sumber kluster sharded: bidang _id di setiap koleksi yang disinkronkan harus unik, atau ketidakkonsistenan data dapat terjadi.

  • Untuk sumber kluster sharded: jumlah node mongos tidak boleh melebihi 10, dan instans tidak boleh berisi dokumen yatim (orphaned documents). Lihat topik FAQ untuk cara menghapus dokumen yatim.

  • Instans ApsaraDB for MongoDB standalone, kluster Azure Cosmos DB for MongoDB, dan kluster elastis Amazon DocumentDB tidak didukung sebagai sumber.

  • DTS tidak dapat terhubung ke MongoDB melalui titik akhir SRV.

  • Oplog harus diaktifkan dan menyimpan data log selama minimal 7 hari, ATAU change streams harus diaktifkan agar DTS dapat berlangganan perubahan data dari 7 hari terakhir. Jika kedua kondisi ini tidak terpenuhi, DTS mungkin gagal mengambil perubahan sumber, dan kehilangan data atau ketidakkonsistenan dapat terjadi di luar cakupan perjanjian tingkat layanan (SLA).

    Penting

    – Gunakan oplog untuk mencatat perubahan database sumber (direkomendasikan).– Change streams memerlukan MongoDB 4.0 atau versi lebih baru. Sinkronisasi dua arah tidak didukung saat menggunakan change streams.– Untuk kluster Amazon DocumentDB non-elastis, aktifkan change streams dan atur Migration Method ke ChangeStream serta Architecture ke Sharded Cluster.

  • Selama sinkronisasi data penuh, jangan memodifikasi skema database atau koleksi, atau mengubah data bertipe ARRAY.

  • Untuk sumber kluster sharded, jangan menjalankan perintah berikut selama sinkronisasi: shardCollection, reshardCollection, unshardCollection, moveCollection, atau movePrimary. Perintah-perintah ini dapat menyebabkan ketidakkonsistenan data.

  • Jika database sumber adalah instans MongoDB yang menggunakan arsitektur kluster sharded dan balancer database sumber melakukan penyeimbangan data, latensi dapat terjadi pada instans tersebut.

Batasan lainnya

  • Hanya sinkronisasi tingkat koleksi yang didukung.

  • Database admin, config, dan local tidak dapat disinkronkan.

  • Rekaman tunggal yang melebihi 10 MB akan menyebabkan tugas gagal.

  • Transaksi tidak dipertahankan. DTS mengonversi setiap transaksi menjadi catatan individual di tujuan.

  • Jika node broker ditambahkan atau dihapus di instans Kafka tujuan saat tugas DTS sedang berjalan, restart tugas DTS tersebut.

  • Pastikan DTS dapat terhubung ke instans sumber dan tujuan. Misalnya, pengaturan keamanan instans database dan parameter listeners serta advertised.listeners dalam file server.properties instans Kafka yang dikelola sendiri tidak membatasi akses dari DTS.

  • Jalankan sinkronisasi selama jam sepi ketika utilisasi CPU pada database sumber dan tujuan berada di bawah 30%.

  • DTS secara otomatis mencoba mengulang instans yang gagal jika telah berjalan kurang dari 7 days. Sebelum mengalihkan lalu lintas ke tujuan, hentikan atau rilis instans sinkronisasi untuk mencegah proses otomatis berlanjut dan menimpa data tujuan.

  • DTS menghitung latensi sinkronisasi data inkremental berdasarkan stempel waktu data tersinkronisasi terbaru di database tujuan dan stempel waktu saat ini di database sumber. Jika tidak ada operasi pembaruan yang dilakukan pada database sumber dalam jangka waktu yang lama, latensi sinkronisasi mungkin tidak akurat. Jika latensi tugas sinkronisasi data terlalu tinggi, Anda dapat melakukan operasi pembaruan pada database sumber untuk memperbarui latensi tersebut.

  • Jika tugas DTS gagal dijalankan, dukungan teknis DTS akan mencoba memulihkan tugas tersebut dalam waktu 8 jam. Selama pemulihan, tugas mungkin direstart, dan parameter tugas mungkin dimodifikasi. Hanya parameter tugas DTS yang mungkin dimodifikasi. Parameter database tidak dimodifikasi. Parameter yang mungkin dimodifikasi mencakup tetapi tidak terbatas pada parameter dalam bagian "Modify instance parameters".

  • Untuk sumber kluster sharded dengan Oplog sebagai metode sinkronisasi inkremental, DTS tidak menjamin urutan penulisan lintas shard ke topik Kafka target.

  • Untuk sumber kluster sharded, nonaktifkan balancer MongoDB selama sinkronisasi data penuh. Aktifkan kembali hanya setelah sinkronisasi penuh selesai dan sinkronisasi inkremental dimulai. Lihat Mengelola balancer ApsaraDB for MongoDB.

Penagihan

Jenis sinkronisasi

Biaya

Sinkronisasi data penuh

Gratis

Sinkronisasi data inkremental

Dikenai biaya. Lihat Ikhtisar penagihan

Jenis sinkronisasi dan operasi yang didukung

Full synchronization menyalin semua data historis dari koleksi MongoDB yang dipilih ke topik Kafka tujuan. DTS mendukung objek DATABASE dan COLLECTION.

Incremental synchronization terus-menerus mengirimkan event perubahan setelah sinkronisasi penuh selesai. Operasi yang didukung bergantung pada metode sinkronisasi inkremental:

Operasi

Oplog

Change streams

INSERT

Didukung

Didukung

UPDATE

Didukung

Didukung

DELETE

Didukung

Didukung

CREATE COLLECTION / INDEX

Didukung

Tidak didukung

DROP DATABASE / COLLECTION / INDEX

Didukung

Hanya DROP DATABASE dan COLLECTION

RENAME COLLECTION

Didukung

Didukung

Sinkronisasi inkremental tidak mengambil database yang dibuat setelah tugas dimulai.

Izin akun database yang diperlukan

Database

Izin yang diperlukan

ApsaraDB for MongoDB sumber

Izin baca pada database yang akan disinkronkan, database admin, dan database local

Untuk instruksi pembuatan akun, lihat Manajemen akun.

Catatan

Jika Anda menggunakan ChangeStream sebagai metode sinkronisasi inkremental, akun database sumber memerlukan izin baca Change Streams di seluruh instans (seperti readAnyDatabase). Jika sumbernya adalah instans ApsaraDB for MongoDB dengan akun kustom, Anda juga harus memberikan izin baca akun tersebut pada database admin. Untuk detailnya, lihat Izin akun root yang ditentukan saat pembuatan instans.

Buat tugas sinkronisasi

Langkah 1: Buka halaman Sinkronisasi Data

Gunakan Konsol DTS atau Konsol DMS.

Konsol DTS

  1. Masuk ke Konsol DTS.Konsol DTS

  2. Di panel navigasi kiri, klik Data Synchronization.

  3. Di pojok kiri atas, pilih wilayah tempat tugas sinkronisasi berada.

Konsol DMS

Langkah-langkah dapat berbeda tergantung pada mode konsol DMS. Lihat Mode simple dan Sesuaikan tata letak dan gaya konsol DMS.
  1. Masuk ke Konsol DMS.Konsol DMS

  2. Di bilah navigasi atas, arahkan kursor ke Data + AI lalu pilih DTS (DTS) > Data Synchronization.

  3. Dari daftar drop-down di sebelah kanan Data Synchronization Tasks, pilih wilayah.

Langkah 2: Konfigurasikan database sumber dan tujuan

Klik Create Task, lalu konfigurasikan parameter yang dijelaskan dalam tabel berikut.

Umum

Parameter

Deskripsi

Task Name

Nama untuk tugas DTS. DTS menghasilkan nama secara otomatis. Tentukan nama deskriptif untuk mengidentifikasi tugas. Nama tidak perlu unik.

Database sumber

Parameter

Deskripsi

Select Existing Connection

Jika instans sumber telah terdaftar di DTS, pilih dari daftar drop-down. DTS akan mengisi kolom lainnya secara otomatis. Jika tidak, konfigurasikan kolom di bawah ini. Di Konsol DMS, pilih dari daftar Select a DMS database instance.

Database Type

Pilih MongoDB.

Access Method

Pilih Alibaba Cloud Instance.

Instance Region

Pilih wilayah instans ApsaraDB for MongoDB sumber.

Replicate Data Across Alibaba Cloud Accounts

Pilih No jika database sumber milik Akun Alibaba Cloud saat ini.

Architecture

Pilih Replica Set untuk contoh ini. Jika sumbernya adalah Sharded Cluster, masukkan juga Shard account dan Shard password.

Migration Method

Metode untuk menyinkronkan data inkremental. Oplog direkomendasikan jika fitur oplog diaktifkan. ChangeStream tersedia jika change streams diaktifkan pada sumber. Lihat Change Streams untuk detailnya. Catatan: Untuk kluster Amazon DocumentDB non-elastis, hanya ChangeStream yang didukung. Saat Architecture diatur ke Sharded Cluster dan Migration Method diatur ke ChangeStream, kolom Shard account dan Shard password tidak diperlukan.

Instance ID

Pilih ID instans instans ApsaraDB for MongoDB sumber.

Authentication Database

Database yang berisi kredensial akun. Default-nya adalah admin.

Database Account

Akun untuk database sumber. Lihat Izin akun database yang diperlukan.

Database Password

Password untuk akun database.

Encryption

Menentukan jenis enkripsi koneksi: Non-encrypted, SSL-encrypted, atau Mongo Atlas SSL. Opsi yang tersedia bergantung pada nilai Access Method dan Architecture. Catatan: SSL-encrypted tidak tersedia saat Architecture adalah Sharded Cluster dan Migration Method adalah Oplog. Untuk MongoDB yang dikelola sendiri dengan arsitektur Replica Set dan SSL-encrypted dipilih, unggah sertifikat CA untuk memverifikasi koneksi.

Database tujuan

Parameter

Deskripsi

Select Existing Connection

Jika instans Kafka tujuan telah terdaftar di DTS, pilih dari daftar drop-down. Jika tidak, konfigurasikan kolom di bawah ini.

Database Type

Pilih Kafka.

Access Method

Pilih Alibaba Cloud Instance.

Instance Region

Pilih wilayah instans Kafka tujuan.

Kafka Instance ID

Pilih ID instans Kafka tujuan.

Encryption

Pilih Non-encrypted atau SCRAM-SHA-256 sesuai kebutuhan keamanan Anda.

Topic

Pilih topik untuk menerima data yang disinkronkan.

Use Kafka Schema Registry

Kafka Schema Registry adalah layanan metadata RESTful untuk menyimpan dan mengambil skema Avro. Pilih No untuk melewati, atau Yespengaturan pemberitahuan peringatan dan berikan URL atau alamat IP Schema Registry.

Langkah 3: Uji konektivitas

Klik Test Connectivity and Proceed di bagian bawah halaman.

DTS secara otomatis menambahkan blok CIDR servernya ke pengaturan keamanan database sumber dan tujuan, jika diizinkan. Untuk pengaturan manual, lihat Tambahkan blok CIDR server DTS.
Untuk database yang dikelola sendiri yang tidak menggunakan Alibaba Cloud Instance sebagai metode akses, klik Test Connectivity di kotak dialog CIDR Blocks of DTS Servers.

Langkah 4: Konfigurasikan objek yang akan disinkronkan

Pada langkah Configure Objects, atur parameter berikut.

Parameter

Deskripsi

Synchronization Types

Incremental Data Synchronization dipilih secara default. Secara opsional pilih juga Full Data Synchronization. Schema Synchronization tidak tersedia. Saat sinkronisasi penuh diaktifkan, DTS terlebih dahulu menyalin semua data historis sebelum memulai sinkronisasi inkremental.

Processing Mode of Conflicting Tables

Precheck and Report Errors: pemeriksaan awal gagal jika tujuan berisi koleksi dengan nama yang sama seperti sumber. Gunakan pemetaan nama objek untuk menyelesaikan konflik penamaan. Ignore Errors and Proceed: melewati pemeriksaan konflik nama. Jika catatan di tujuan memiliki kunci primer atau kunci unik yang sama dengan catatan sumber, catatan tujuan dipertahankan.

Peringatan

Opsi ini dapat menyebabkan ketidakkonsistenan data.

Data Format in Kafka

Hanya Canal JSON yang didukung.

Kafka Data Compression Format

Format kompresi untuk data yang ditulis ke Kafka. Opsi: LZ4 (default — rasio kompresi rendah, kecepatan tinggi), GZIP (rasio kompresi tinggi, kecepatan rendah, penggunaan CPU tinggi), Snappy (seimbang).

Policy for Shipping Data to Kafka Partitions

Pilih kebijakan perutean partisi sesuai kebutuhan Anda.

Message acknowledgement mechanism

Pilih mekanisme pengakuan pesan sesuai kebutuhan Anda.

Topic That Stores DDL Information

Pilih topik untuk menyimpan event DDL. Jika dibiarkan kosong, event DDL ditulis ke topik yang sama dengan event data.

Capitalization of Object Names in Destination Instance

Mengontrol kapitalisasi nama database dan koleksi di tujuan. Default-nya adalah DTS default policy. Lihat Tentukan kapitalisasi nama objek di instans tujuan.

Source Objects

Pilih objek dari bagian Source Objects dan klik ikon panah untuk memindahkannya ke Selected Objects. Hanya pemilihan tingkat koleksi yang didukung.

Langkah 5: Konfigurasikan pengaturan lanjutan

Klik Next: Advanced Settings dan konfigurasikan parameter berikut.

Parameter

Deskripsi

Dedicated Cluster for Task Scheduling

Secara default, DTS menjadwalkan tugas pada kluster bersama. Beli klaster khusus untuk stabilitas yang lebih tinggi.

Retry Time for Failed Connections

Berapa lama DTS mencoba ulang saat sumber atau tujuan tidak dapat dijangkau. Rentang: 10–1440 menit. Default: 720 menit. Kami merekomendasikan Anda mengatur parameter ini ke nilai lebih dari 30. Jika Anda menentukan waktu coba ulang berbeda untuk beberapa tugas yang berbagi database yang sama, nilai terpendek yang berlaku. DTS mengenakan biaya untuk waktu coba ulang.

Retry Time for Other Issues

Berapa lama DTS mencoba ulang saat operasi DDL atau DML gagal. Rentang: 1–1440 menit. Default: 10 menit. Atur minimal 10 menit. Harus lebih pendek dari Retry Time for Failed Connections.

Obtain the entire document after it is updated

Hanya tersedia saat Migration Method adalah ChangeStream. Yes: menyinkronkan dokumen lengkap untuk setiap pembaruan, yang dapat meningkatkan beban sumber dan menyebabkan latensi. Jika DTS tidak dapat mengambil dokumen lengkap, hanya bidang yang diperbarui yang dikirim. No: hanya menyinkronkan bidang yang berubah.

Enable Throttling for Full Data Synchronization

Saat diaktifkan, konfigurasikan QPS to the source database, RPS of Full Data Migration, dan Data migration speed for full migration (MB/s) untuk mengurangi beban pada tujuan. Hanya tersedia saat Full Data Synchronization dipilih.

Only one data type for primary key _id in a table of the data to be synchronized

Hanya tersedia saat Full Data Synchronization dipilih. Yes: DTS melewati pemindaian tipe data bidang _id selama sinkronisasi penuh. No: DTS memindai tipe data bidang _id.

Enable Throttling for Incremental Data Synchronization

Saat diaktifkan, konfigurasikan RPS of Incremental Data Synchronization dan Data synchronization speed for incremental synchronization (MB/s) untuk mengurangi beban pada tujuan.

Environment Tag

Tag opsional untuk mengidentifikasi instans.

Configure ETL

Aktifkan untuk menerapkan logika ekstrak, transformasi, muat (ETL). Masukkan pernyataan pemrosesan data di editor kode. Lihat Konfigurasikan ETL dalam tugas migrasi data atau sinkronisasi data.

Monitoring and Alerting

Saat diaktifkan, tentukan ambang batas peringatan dan pengaturan notifikasi. DTS memberi tahu kontak peringatan saat tugas gagal atau latensi sinkronisasi melebihi ambang batas. Lihat Konfigurasikan pemantauan dan peringatan.

Langkah 6: Jalankan pemeriksaan awal

Klik Next: Save Task Settings and Precheck.

Untuk melihat pratinjau parameter API untuk konfigurasi ini, arahkan kursor ke tombol tersebut lalu klik Preview OpenAPI parameters sebelum melanjutkan.
Tugas tidak dapat dimulai sampai pemeriksaan awal berhasil.
Jika pemeriksaan awal gagal, klik View Details di sebelah setiap item yang gagal, selesaikan masalahnya, lalu klik Precheck Again.
Jika muncul peringatan untuk suatu item: perbaiki peringatan yang tidak dapat diabaikan sebelum melanjutkan. Untuk peringatan yang dapat diabaikan, klik Confirm Alert Details, lalu Ignore di dialog, konfirmasi, dan klik Precheck Again. Mengabaikan peringatan dapat menyebabkan ketidakkonsistenan data.

Langkah 7: Beli dan mulai instans

  1. Tunggu hingga Success Rate mencapai 100%, lalu klik Next: Purchase Instance.

  2. Di halaman buy, konfigurasikan parameter berikut.

    Parameter

    Deskripsi

    Billing Method

    Subscription: bayar di muka; lebih hemat biaya untuk penggunaan jangka panjang. Pay-as-you-go: ditagih per jam; cocok untuk penggunaan jangka pendek. Lepaskan instans saat tidak lagi diperlukan untuk menghentikan penagihan.

    Resource Group Settings

    Kelompok sumber daya untuk instans sinkronisasi. Default: default resource group. Lihat Apa itu Resource Management?

    Instance Class

    Tingkat kecepatan sinkronisasi. Lihat Kelas instans instansi sinkronisasi data.

    Subscription Duration

    Tersedia untuk penagihan Subscription. Opsi: 1–9 bulan, 1 tahun, 2 tahun, 3 tahun, atau 5 tahun.

  3. Baca dan terima Data Transmission Service (Pay-as-you-go) Service Terms.

  4. Klik Buy and Start, lalu klik OK di dialog konfirmasi.

Tugas muncul di daftar tugas. Pantau perkembangannya di sana.

Konfigurasikan pemetaan koleksi ke topik

Secara default, semua koleksi ditulis ke topik yang dipilih dalam konfigurasi database tujuan. Untuk mengarahkan koleksi tertentu ke topik berbeda:

  1. Di area Selected Objects, arahkan kursor ke nama topik tujuan di tingkat koleksi.

  2. Klik Edit di sebelah nama topik.

  3. Di dialog Edit Table, konfigurasikan pengaturan berikut.

    Parameter

    Deskripsi

    Name of target Topic

    Topik tujuan untuk koleksi ini. Topik harus ada di instans Kafka. Jika dimodifikasi, data ditulis ke topik baru tersebut. Default: topik yang dipilih selama konfigurasi database tujuan.

    Filter Conditions

    Filter baris opsional. Lihat Filter data tugas dengan kondisi SQL.

    Number of Partitions

    Jumlah partisi untuk menulis data ke topik.

  4. Klik OK.

Contoh pengiriman data

Setiap perubahan inkremental dari MongoDB diserialisasi sebagai pesan Canal JSON dan dikirimkan ke topik Kafka yang dikonfigurasi. Struktur pesan bervariasi tergantung pada metode sinkronisasi inkremental dan pengaturan pembaruan.

Pilih skenario Anda

Tujuan

Konfigurasi

Sinkronisasi latensi rendah dengan dukungan DDL lengkap

Atur Migration Method ke Oplog (Skenario 1)

Change Streams dengan pembaruan dokumen parsial

Atur Migration Method ke ChangeStream dan Obtain the entire document after it is updated ke No (Skenario 2)

Change Streams dengan dokumen lengkap pada setiap pembaruan

Atur Migration Method ke ChangeStream dan Obtain the entire document after it is updated ke Yes (Skenario 3)

Bidang pesan Canal JSON

Ketiga skenario menggunakan envelope Canal JSON tingkat atas yang sama. Bidang berikut muncul di setiap pesan.

Bidang

Tipe

Deskripsi

database

string

Nama database sumber

table

string

Nama koleksi sumber

type

string

Jenis operasi: INSERT, UPDATE, DELETE, atau DDL

isDdl

boolean

true untuk event DDL (drop koleksi, rename); false untuk event DML

es

number

Stempel waktu event di database sumber (milidetik Unix)

ts

number

Stempel waktu saat DTS memproses event (milidetik Unix)

id

number

ID event internal DTS

pkNames

array

Nama bidang kunci primer (biasanya ["_id"])

data

array

Data dokumen setelah operasi. Untuk pembaruan parsial (Oplog atau ChangeStream tanpa pengambilan dokumen lengkap), hanya berisi bidang yang berubah menggunakan operator pembaruan MongoDB ($set, $unset). null untuk event DDL.

old

array

Status dokumen sebelum pembaruan: hanya bidang _id yang disertakan. null untuk event INSERT dan DDL.

sql

object atau null

Detail pernyataan DDL (untuk isDdl: true). null untuk event DML.

gtid

null

Tidak berlaku untuk MongoDB (disediakan untuk kompatibilitas MySQL).

mysqlType

null

Tidak berlaku untuk MongoDB.

serverId

null

Tidak berlaku untuk MongoDB.

sqlType

null

Tidak berlaku untuk MongoDB.

Skenario 1: Oplog

Atur Migration Method ke Oplog.

Jenis perubahan sumber

Pernyataan sumber

Data yang diterima oleh topik tujuan

insert

db.kafka_test.insert({"cid":"a","person":{"name":"testName","age":NumberInt(18),"skills":["database","ai"]}})

Lihat contoh di bawah

update $set

db.kafka_test.update({"cid":"a"},{$set:{"person.age":NumberInt(20)}})

Lihat contoh di bawah

update $set new field

db.kafka_test.update({"cid":"a"},{$set:{"salary":100}})

Lihat contoh di bawah

update $unset (remove field)

db.kafka_test.update({"cid":"a"},{$unset:{"salary":1}})

Lihat contoh di bawah

delete

db.kafka_test.deleteOne({"cid":"a"})

Lihat contoh di bawah

ddl drop

db.kafka_test.drop()

Lihat contoh di bawah

Lihat data (klik untuk membuka)

{
    "data": [{
        "person": {
            "skills": ["database", "ai"],
            "name": "testName",
            "age": 18
        },
        "_id": {
            "$oid": "67d27da49591697476e1****"
        },
        "cid": "a"
    }],
    "database": "kafkadb",
    "es": 1741847972000,
    "gtid": null,
    "id": 174184797200000****,
    "isDdl": false,
    "mysqlType": null,
    "old": null,
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741847973438,
    "type": "INSERT"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "$set": {
            "person.age": 20
        }
    }],
    "database": "kafkadb",
    "es": 1741848051000,
    "gtid": null,
    "id": 174184805100000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848051984,
    "type": "UPDATE"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "$set": {
            "salary": 100.0
        }
    }],
    "database": "kafkadb",
    "es": 1741848146000,
    "gtid": null,
    "id": 174184814600000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848147734,
    "type": "UPDATE"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "$unset": {
            "salary": true
        }
    }],
    "database": "kafkadb",
    "es": 1741848207000,
    "gtid": null,
    "id": 174184820700000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848208186,
    "type": "UPDATE"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "database": "kafkadb",
    "es": 1741848289000,
    "gtid": null,
    "id": 174184828900000****,
    "isDdl": false,
    "mysqlType": null,
    "old": null,
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848289798,
    "type": "DELETE"
}

Lihat data (klik untuk membuka)

{
    "data": null,
    "database": "kafkadb",
    "es": 1741847893000,
    "gtid": null,
    "id": 1741847893000000005,
    "isDdl": true,
    "mysqlType": null,
    "old": null,
    "pkNames": null,
    "serverId": null,
    "sql": {
        "drop": "kafka_test"
    },
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741847893760,
    "type": "DDL"
}

Skenario 2: ChangeStream — hanya bidang yang diperbarui

Atur Migration Method ke ChangeStream. Atur Obtain the entire document after it is updated ke No.

Pesan insert dan delete identik dengan Skenario 1. Pesan update hanya berisi bidang yang berubah.

Lihat data (klik untuk membuka)

{
    "data": [{
        "$set": {
            "person.age": 20
        }
    }],
    "database": "kafkadb",
    "es": 1741848051000,
    "gtid": null,
    "id": 174184805100000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848052912,
    "type": "UPDATE"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "$unset": {
            "salary": 1
        }
    }],
    "database": "kafkadb",
    "es": 1741848207000,
    "gtid": null,
    "id": 174184820700000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848209142,
    "type": "UPDATE"
}

Skenario 3: ChangeStream — dokumen lengkap saat pembaruan

Atur Migration Method ke ChangeStream. Atur Obtain the entire document after it is updated ke Yes.

Event update mengirimkan dokumen lengkap setelah perubahan, bukan hanya bidang yang berubah.

Lihat data (klik untuk membuka)

{
    "data": [{
        "person": {
            "skills": ["database", "ai"],
            "name": "testName",
            "age": 20
        },
        "_id": {
            "$oid": "67d27da49591697476e1****"
        },
        "cid": "a"
    }],
    "database": "kafkadb",
    "es": 1741848051000,
    "gtid": null,
    "id": 174184805100000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848052219,
    "type": "UPDATE"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "person": {
            "skills": ["database", "ai"],
            "name": "testName",
            "age": 20
        },
        "_id": {
            "$oid": "67d27da49591697476e1****"
        },
        "salary": 100.0,
        "cid": "a"
    }],
    "database": "kafkadb",
    "es": 1741848146000,
    "gtid": null,
    "id": 174184814600000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848147327,
    "type": "UPDATE"
}

Lihat data (klik untuk membuka)

{
    "data": [{
        "person": {
            "skills": ["database", "ai"],
            "name": "testName",
            "age": 20
        },
        "_id": {
            "$oid": "67d27da49591697476e1****"
        },
        "cid": "a"
    }],
    "database": "kafkadb",
    "es": 1741848207000,
    "gtid": null,
    "id": 174184820700000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "_id": {
            "$oid": "67d27da49591697476e1****"
        }
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "kafka_test",
    "ts": 1741848208401,
    "type": "UPDATE"
}

Kasus khusus: fullDocument hilang di ChangeStream

Saat bidang fullDocument pada event pembaruan ChangeStream hilang — misalnya, saat dokumen berpindah lintas shard dalam koleksi sharded — pesan yang dikirimkan kembali ke perilaku Oplog (pembaruan parsial dengan operator $set atau $unset).

Contoh: pembaruan koleksi sharded di mana fullDocument hilang

Data dasar sumber:

use admin
db.runCommand({ enablesharding:"dts_test" })

Perubahan inkremental sumber:

use dts_test
sh.shardCollection("dts_test.cstest",{"name":"hashed"})
db.cstest.insert({"_id":1,"name":"a"})
db.cstest.updateOne({"_id":1,"name":"a"},{$set:{"name":"b"}})

Lihat data (klik untuk membuka)

{
    "data": [{
        "$set": {
            "name": "b"
        }
    }],
    "database": "dts_test",
    "es": 1740720994000,
    "gtid": null,
    "id": 174072099400000****,
    "isDdl": false,
    "mysqlType": null,
    "old": [{
        "name": "a",
        "_id": 1.0
    }],
    "pkNames": ["_id"],
    "serverId": null,
    "sql": null,
    "sqlType": null,
    "table": "cstest",
    "ts": 1740721007099,
    "type": "UPDATE"
}

FAQ

Dapatkah saya mengubah format kompresi data Kafka setelah tugas dibuat?

Ya. Lihat Modify the objects to be synchronized.

Apakah saya dapat mengubah Message acknowledgment mechanism setelah tugas dibuat?

Ya. Lihat Modify the objects to be synchronized.

Langkah selanjutnya