Setelah membuat tugas pelacakan perubahan, Anda dapat menggunakan kit pengembangan perangkat lunak (software development kit/SDK) yang disediakan oleh Data Transmission Service (DTS) untuk berlangganan perubahan data. Topik ini menjelaskan cara menggunakan SDK untuk mengonsumsi data dari sumber data terdistribusi, seperti PolarDB-X 1.0 dan database logis DMS.
Prasyarat
Anda telah membuat instans pelacakan perubahan yang berada dalam status Normal. Untuk informasi selengkapnya, lihat Buat tugas pelacakan perubahan untuk instans PolarDB-X 1.0 atau Buat tugas pelacakan perubahan untuk database logis DMS.
-
Anda telah membuat kelompok konsumen untuk instans langganan Anda.
-
Jika Anda menggunakan RAM user untuk mengonsumsi data yang berlangganan, RAM user tersebut harus memiliki izin AliyunDTSFullAccess dan izin akses ke objek yang berlangganan. Untuk informasi selengkapnya, lihat Berikan izin kepada RAM user untuk mengelola DTS menggunakan kebijakan sistem dan Kelola izin RAM user.
Catatan penggunaan
-
Saat mengonsumsi data yang berlangganan, Anda harus memanggil metode commit dari DefaultUserRecord untuk melakukan commit informasi offset. Jika tidak, data mungkin dikonsumsi berulang kali.
-
Proses konsumsi yang berbeda saling independen satu sama lain.
Prosedur
Unduh dan ekstrak kode contoh SDK.
-
Verifikasi versi kode SDK.
-
Buka direktori tempat Anda mengekstrak kode contoh SDK.
-
Gunakan editor teks untuk membuka file pom.xml di direktori tersebut.
-
Perbarui SDK pelacakan perubahan ke versi terbaru.
CatatanAnda dapat menemukan dependensi Maven terbaru di halaman dts-new-subscribe-sdk.
-
Edit kode SDK.
Anda dapat membuka file yang telah didekompresi menggunakan perangkat lunak pengodean.
Berdasarkan pola penggunaan klien SDK, buka file DistributedDTSConsumerDemo.java.
CatatanJalur file Java adalah
aliyun-dts-subscribe-sdk-java-master/src/test/java/com/aliyun/dts/subscribe/clients/.Atur parameter dalam kode Java.
public static void main(String[] args) throws ClientException { // Konfigurasi untuk berlangganan ke sumber data terdistribusi, seperti PolarDB-X 1.0 (sebelumnya DRDS). Konfigurasikan informasi seperti AccessKey, ID instans, ID tugas utama, dan kelompok konsumen. String accessKeyId = "LTA***********99reZ"; String accessKeySecret = "****************"; String regionId = "cn-hangzhou"; String dtsInstanceId = "dtse5212sed162****"; String jobId = "l791216x16d****"; String sid = "dtsip412t13160****"; String userName = "xftest"; String password = "******"; String proxyUrl = "dts-cn-****.com:18001"; // checkpoint awal untuk pencarian pertama (stempel waktu yang ditetapkan, misalnya 1566180200 jika Anda menginginkan (Mon Aug 19 10:03:21 CST 2019)) String checkpoint = "1639620090"; // Konversi nama database/tabel fisik ke nama database/tabel logis boolean mapping = true; // jika memaksa menggunakan checkpoint konfigurasi saat memulai. Untuk reset checkpoint, hanya mode assign yang berfungsi boolean isForceUseInitCheckpoint = false; ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.ASSIGN; DistributedDTSConsumerDemo demo = new DistributedDTSConsumerDemo(userName, password, regionId, jobId, sid, dtsInstanceId, accessKeyId, accessKeySecret, subscribeMode, proxyUrl, checkpoint, isForceUseInitCheckpoint, mapping); demo.start(); }Parameter
Deskripsi
Cara memperoleh
accessKeyId
ID AccessKey.
Untuk informasi selengkapnya, lihat Peroleh Pasangan AccessKey.
accessKeySecret
Rahasia AccessKey.
regionId
ID wilayah tempat tugas pelacakan perubahan berada.
Di Konsol DTS, klik ID instans pelacakan perubahan yang dituju. Di halaman Basic Information, Anda dapat mengambil informasi wilayah. Misalnya, jika wilayahnya adalah Tiongkok (Hangzhou), atur parameter ini menjadi
cn-hangzhou. Untuk informasi selengkapnya, lihat Daftar wilayah.dtsInstanceId
ID instans pelacakan perubahan.
Di Konsol DTS, klik ID instans pelacakan perubahan yang dituju. Di halaman Basic Information, Anda dapat mengambil DTS Instance ID dari instans pelacakan perubahan tersebut.
jobId
ID tugas pelacakan perubahan.
Anda dapat memanggil operasi DescribeDtsJobs untuk mengambil ID tugas pelacakan perubahan (DtsJobId).
sid
ID kelompok konsumen.
Di Konsol DTS, klik ID instans pelacakan perubahan yang dituju. Di panel navigasi sebelah kiri, klik Consume Data. Anda dapat mengambil Consumer Group ID/Name dan Account dari kelompok konsumen tersebut.
CatatanKata sandi akun kelompok konsumen ditentukan saat Anda membuat kelompok konsumen.
userName
Akun kelompok konsumen.
password
Kata sandi akun kelompok konsumen.
proxyUrl
Titik akhir dan port saluran pelacakan perubahan.
CatatanJika instance ECS tempat Anda men-deploy klien SDK dan saluran pelacakan perubahan berada dalam jaringan klasik atau virtual private cloud (VPC) yang sama, berlanggananlah data melalui jaringan internal untuk mencapai latensi terendah.
-
Kami tidak merekomendasikan penggunaan titik akhir publik karena potensi ketidakstabilan jaringan.
Di Konsol DTS, klik ID instans pelacakan perubahan yang dituju. Di halaman Basic Information, Anda dapat mengambil informasi Network.
checkpoint
Offset konsumen. Ini adalah stempel waktu mulai SDK client mengonsumsi catatan data. Nilainya merupakan Stempel waktu UNIX dalam satuan detik.
CatatanAnda dapat menggunakan informasi offset konsumen untuk hal-hal berikut:
Jika proses konsumsi terganggu, Anda dapat meneruskan offset konsumen untuk melanjutkan konsumsi data dan mencegah kehilangan data.
Saat memulai klien SDK, Anda dapat meneruskan offset konsumen yang diperlukan untuk menyesuaikan offset langganan dan mengonsumsi data sesuai kebutuhan.
Offset konsumen harus berada dalam rentang waktu instans pelacakan perubahan dan harus dikonversi ke Stempel waktu UNIX.
CatatanAnda dapat melihat rentang waktu instans pelacakan perubahan di kolom Data Range pada daftar tugas pelacakan.
Anda dapat menggunakan mesin pencari untuk menemukan konverter Stempel waktu UNIX.
Opsi: Untuk memodifikasi tipe data dari data yang berlangganan, Anda dapat memodifikasi metode
buildRecordListener()atau menggunakan kelas kustom.public static Map<String, RecordListener> buildRecordListener() { // pengguna dapat mengimplementasikan listener sendiri RecordListener mysqlRecordPrintListener = new RecordListener() { @Override public void consume(DefaultUserRecord record) { OperationType operationType = record.getOperationType(); if (operationType.equals(OperationType.INSERT) || operationType.equals(OperationType.UPDATE) || operationType.equals(OperationType.DELETE) || operationType.equals(OperationType.HEARTBEAT)) { // konsumsi catatan RecordListener recordPrintListener = new DefaultRecordPrintListener(DbType.MySQL); recordPrintListener.consume(record); // metode commit mendorong pembaruan checkpoint record.commit(""); } } }; return Collections.singletonMap("mysqlRecordPrinter", mysqlRecordPrintListener); }Buka struktur proyek di IDE Anda dan pastikan versi OpenJDK untuk proyek adalah 1.8.
Jalankan kode klien.
Output menunjukkan bahwa klien sedang berlangganan perubahan data dari database sumber.
Klien SDK secara berkala mengumpulkan dan menampilkan statistik tentang konsumsi data. Statistik tersebut mencakup jumlah total catatan data yang dikirim dan diterima, volume data total, serta catatan per detik (records per second/RPS).
Tabel 1. Statistik konsumsi data
Parameter
Deskripsi
outCountsJumlah total catatan data yang dikonsumsi oleh klien SDK.
outBytesVolume total data yang dikonsumsi oleh klien SDK, dalam byte.
outRpsJumlah permintaan per detik yang dikirim oleh klien SDK untuk mengonsumsi data.
outBpsJumlah bit yang ditransmisikan per detik saat klien SDK mengonsumsi data.
countTidak ada.
inBytesVolume total data yang dikirim oleh server DTS, dalam byte.
DStoreRecordQueueUkuran antrian cache data saat server DTS mengirim data.
inCountsJumlah total catatan data yang dikirim oleh server DTS.
inRpsJumlah permintaan yang dikirim oleh server DTS per detik.
inBpsJumlah bit yang ditransmisikan per detik saat server DTS mengirim data.
__dtStempel waktu saat klien SDK menerima data, dalam milidetik.
DefaultUserRecordQueueUkuran antrian cache data setelah serialisasi.