All Products
Search
Document Center

Data Transmission Service:Gunakan SDK untuk mengonsumsi data yang dilacak dari instans PolarDB-X 1.0

Last Updated:Aug 22, 2026

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

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

  1. Unduh dan ekstrak kode contoh SDK.

  2. Verifikasi versi kode SDK.

    1. Buka direktori tempat Anda mengekstrak kode contoh SDK.

    2. Gunakan editor teks untuk membuka file pom.xml di direktori tersebut.

    3. Perbarui SDK pelacakan perubahan ke versi terbaru.

      Catatan

      Anda dapat menemukan dependensi Maven terbaru di halaman dts-new-subscribe-sdk.

      Lokasi parameter versi SDK (klik untuk memperluas)

      <name>dts-new-subscribe-sdk</name>
      <url>https://www.aliyun.com/product/dts</url>
      <description>The Aliyun new Subscribe SDK for Java used for accessing Data Transmission Service</description>
      <packaging>jar</packaging>
      <groupId>com.aliyun.dts</groupId>
      <artifactId>dts-new-subscribe-sdk</artifactId>
      <version>2.1.4</version>
  3. Edit kode SDK.

    1. Anda dapat membuka file yang telah didekompresi menggunakan perangkat lunak pengodean.

    2. Berdasarkan pola penggunaan klien SDK, buka file DistributedDTSConsumerDemo.java.

      Catatan

      Jalur file Java adalah aliyun-dts-subscribe-sdk-java-master/src/test/java/com/aliyun/dts/subscribe/clients/.

    3. 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.

      Catatan

      Kata 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.

      Catatan
      • Jika 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.

      Catatan

      Anda 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.

      Catatan
      • Anda 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.

  4. 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);
        }
  5. Buka struktur proyek di IDE Anda dan pastikan versi OpenJDK untuk proyek adalah 1.8.

  6. 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

      outCounts

      Jumlah total catatan data yang dikonsumsi oleh klien SDK.

      outBytes

      Volume total data yang dikonsumsi oleh klien SDK, dalam byte.

      outRps

      Jumlah permintaan per detik yang dikirim oleh klien SDK untuk mengonsumsi data.

      outBps

      Jumlah bit yang ditransmisikan per detik saat klien SDK mengonsumsi data.

      count

      Tidak ada.

      inBytes

      Volume total data yang dikirim oleh server DTS, dalam byte.

      DStoreRecordQueue

      Ukuran antrian cache data saat server DTS mengirim data.

      inCounts

      Jumlah total catatan data yang dikirim oleh server DTS.

      inRps

      Jumlah permintaan yang dikirim oleh server DTS per detik.

      inBps

      Jumlah bit yang ditransmisikan per detik saat server DTS mengirim data.

      __dt

      Stempel waktu saat klien SDK menerima data, dalam milidetik.

      DefaultUserRecordQueue

      Ukuran antrian cache data setelah serialisasi.