Setelah mengonfigurasi instans Data Subscription, gunakan kode contoh SDK yang disediakan oleh Data Transmission Service (DTS) untuk mengonsumsi data perubahan.
Prosedur
-
Jika sumber data adalah instans PolarDB-X 1.0 atau database logis DMS, lihat Mengonsumsi data Langganan Data dari PolarDB-X 1.0 menggunakan kode contoh SDK.
-
Jika Anda menggunakan pengguna RAM untuk mengonsumsi data, pengguna RAM tersebut harus memiliki izin AliyunDTSFullAccess serta izin untuk mengakses objek langganan. Untuk informasi selengkapnya tentang cara memberikan izin, lihat Memberikan otorisasi kepada pengguna RAM untuk mengelola instans DTS menggunakan kebijakan sistem dan Mengelola izin pengguna RAM.
-
Setiap konsumen beroperasi secara independen.
-
Topik ini menyediakan klien SDK contoh dalam Java. Untuk kode contoh dalam Python dan Go, lihat dts-subscribe-demo.
Prosedur berikut menjelaskan cara menjalankan kode contoh SDK untuk mengonsumsi data Langganan Data di IntelliJ IDEA (Community Edition 2020.1 untuk Windows).
-
Buat instans Data Subscription. Untuk informasi selengkapnya, lihat Membuat saluran Langganan Data untuk instans ApsaraDB RDS for MySQL, Membuat saluran Langganan Data untuk kluster PolarDB for MySQL, atau Membuat saluran Langganan Data untuk database Oracle.
-
Buat satu atau beberapa kelompok konsumen. Untuk informasi selengkapnya, lihat Membuat kelompok konsumen.
PentingSaat mengonsumsi data Langganan Data, Anda harus memanggil metode
commitdariDefaultUserRecorduntuk melakukan commit checkpoint. Jika tidak, hal ini dapat menyebabkan konsumsi data duplikat. -
Gunakan kode contoh SDK sesuai kebutuhan bisnis Anda.
-
Gunakan paket SDK Langganan Data versi baru (disarankan)
-
Buka IntelliJ IDEA, lalu klik Create New Project untuk membuat proyek bagi aplikasi Anda.
-
Dalam proyek tersebut, temukan file model objek proyek (POM): pom.xml.
-
Tambahkan dependensi berikut ke file pom.xml:
<dependency> <groupId>com.aliyun.dts</groupId> <artifactId>dts-new-subscribe-sdk</artifactId> <version>{dts_new_sdk_version}</version> </dependency>CatatanAnda dapat menemukan dependensi Maven terbaru di halaman dts-new-subscribe-sdk.
-
Untuk informasi selengkapnya tentang cara menggunakan SDK langganan versi baru, lihat Gunakan kode contoh.
-
-
Gunakan versi kustom SDK Langganan Data versi baru
-
Unduh paket kode contoh SDK, lalu ekstrak.
CatatanKlik
dan pilih Download ZIP untuk mengunduh paket tersebut. -
Buka direktori kode contoh SDK yang telah diekstrak. Gunakan editor teks untuk membuka file pom.xml dan perbarui SDK Langganan Data ke versi terbaru.
<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>1.3</version>PentingAnda dapat memperoleh versi terbaru SDK Langganan Data dari situs web Maven. Untuk informasi selengkapnya, lihat halaman Maven untuk SDK Langganan Data.
-
Buka IntelliJ IDEA. Di layar sambutan, klik Open or Import.
-
Pada kotak dialog Open File or Project, arahkan ke direktori kode contoh SDK yang telah diekstrak, pilih file pom.xml, lalu klik OK.
-
Pada kotak dialog yang muncul, pilih Open as Project.
-
Di panel Project IntelliJ IDEA, buka
src > test > java > com.aliyun.dts.subscribe.clients. File demo contoh terdiri dari DTSConsumerAssignDemo dan DTSConsumerSubscribeDemo. Berdasarkan mode penggunaan klien SDK, pilih dan klik ganda file Java yang sesuai: DTSConsumerAssignDemo.java atau DTSConsumerSubscribeDemo.java.Pohon file proyek:
aliyun-dts-subscribe-sdk-java-master [dts-new-subscribe-sdk] D:\aliyun-dts-sub .idea src main test java com.aliyun.dts.subscribe.clients DTSConsumerAssignDemo DTSConsumerSubscribeDemo UserMetaStore target .gitignore dts-new-subscribe-sdk.iml LICENSE pom.xml README.md External Libraries Scratches and ConsolesCatatanDTS mendukung mode penggunaan klien SDK berikut:
-
Mode ASSIGN: Untuk memastikan urutan global pesan, DTS hanya menetapkan satu partisi (partisi 0) untuk setiap topik langganan. Saat menggunakan klien SDK dalam mode ASSIGN, kami menyarankan hanya menjalankan satu klien.
-
Mode SUBSCRIBE: Untuk memastikan urutan global pesan, DTS hanya menetapkan satu partisi (partisi 0) untuk setiap topik langganan. Jika Anda menggunakan klien SDK dalam mode SUBSCRIBE, Anda dapat menjalankan beberapa klien SDK dalam satu kelompok konsumen untuk pemulihan bencana. Jika klien aktif gagal, klien SDK lain akan secara otomatis ditetapkan ke partisi 0 untuk melanjutkan konsumsi.
-
-
-
-
Atur parameter yang diperlukan dalam file Java.
public static void main(String[] args) { // kafka broker url String brokerUrl = "dts-cn-xxx18001"; // topic to consume, partition is 0 String topic = "cn_hangzhou_rm_xxx"; // user password and sid for auth String sid = "dtsxxx"; String userName = "dtstest"; String password = "xxx"; // initial checkpoint for first seek(a timestamp to set, eg 1566180200 if you want (Mon Aug 19 10:03:21 CST 2019)) String initCheckpoint = "1620811813"; // when use subscribe mode, group config is required. kafka consumer group is enabled ConsumerContext.ConsumerSubscribeMode subscribeMode = ConsumerContext.ConsumerSubscribeMode.ASSIGN; // if force use config checkpoint when start. for checkpoint reset, only assign mode works boolean isForceUseInitCheckpoint = true;Tabel 1. Parameter yang diperlukan
Parameter
Deskripsi
Sumber
brokerUrlTitik akhir dan nomor port instans Langganan Data.
Catatan-
Jika instans ECS yang menjalankan klien SDK dan instans Langganan Data berada dalam jaringan klasik atau Virtual Private Cloud (VPC) yang sama, kami menyarankan Anda menggunakan titik akhir internal untuk langganan guna meminimalkan latensi jaringan.
-
Kami tidak menyarankan penggunaan titik akhir publik karena potensi ketidakstabilan jaringan.
Di Konsol DTS, klik ID instans Langganan Data yang dituju. Di halaman Basic Information, Anda dapat memperoleh titik akhir dan nomor port di bagian Network.
topicTopik langganan untuk instans tersebut.
Di Konsol DTS, klik ID instans Langganan Data yang dituju. Di halaman Basic Information, Anda dapat memperoleh Topic di bagian Basic Information.
sidID kelompok konsumen.
Di Konsol DTS, klik ID instans Langganan Data yang dituju, lalu klik Consume Data. Anda dapat memperoleh Consumer Group ID dan Account kelompok konsumen.
CatatanKata sandi untuk username kelompok konsumen ditentukan saat Anda membuat kelompok konsumen.
userNameUsername untuk kelompok konsumen.
PeringatanJika Anda tidak menggunakan klien yang disediakan dalam topik ini, Anda harus mengatur username dalam format
<Username>-<Consumer Group ID>. Contoh:dtstest-dtsae******bpv. Jika tidak, koneksi akan gagal.passwordKata sandi untuk username tersebut.
initCheckpointTitik pemeriksaan konsumsi, yang ditentukan sebagai stempel waktu UNIX, tempat klien SDK mulai mengonsumsi data. Contoh: 1620962769.
CatatanAnda dapat menggunakan informasi titik pemeriksaan konsumsi dalam skenario berikut:
-
Untuk melanjutkan konsumsi dan mencegah kehilangan data setelah gangguan aplikasi, berikan titik pemeriksaan konsumsi terakhir yang diketahui.
-
Saat memulai klien, Anda dapat memberikan titik pemeriksaan konsumsi tertentu untuk mengonsumsi data dari posisi yang diinginkan.
Titik pemeriksaan konsumsi harus berada dalam rentang data instans Langganan Data (seperti yang ditunjukkan pada gambar) dan harus dikonversi ke stempel waktu UNIX.
CatatanAnda dapat menggunakan mesin pencari untuk menemukan konverter stempel waktu UNIX.
ConsumerContext.ConsumerSubscribeMode subscribeModeMode penggunaan klien SDK. Nilai yang valid:
-
ConsumerContext.ConsumerSubscribeMode.ASSIGN: Mode ASSIGN. Hanya satu klien SDK dalam satu kelompok konsumen yang dapat mengonsumsi data Langganan Data. -
ConsumerContext.ConsumerSubscribeMode.SUBSCRIBE: Mode SUBSCRIBE. Anda dapat menjalankan beberapa klien SDK dalam kelompok konsumen yang sama untuk disaster recovery.
N/A
-
-
Di bilah navigasi atas IntelliJ IDEA, pilih untuk menjalankan klien.
CatatanPertama kali menjalankan klien, mungkin memerlukan waktu untuk memuat dan menginstal dependensi yang diperlukan secara otomatis.
-
Setelah dieksekusi, klien SDK berhasil mengonsumsi data perubahan dari database sumber. Berikut adalah contoh catatan UPDATE yang dikonsumsi:
[2021-05-18 16:49:50,260] INFO RecordID [559686] RecordTimestamp [1621327772] Source [{"sourceType": "MySQL", "version": "8.0.18"}] RecordType [UPDATE] Schema info [{, recordFields= [{fieldName='orderid', rawDataTypeNum=3, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}, {fieldName='username', rawDataTypeNum=254, isPrimaryKey=false, isUniqueKey=false, fieldPosition=1}, {fieldName='ordertime', rawDataTypeNum=12, isPrimaryKey=false, isUniqueKey=false, fieldPosition=2}, {fieldName='commodity', rawDataTypeNum=253, isPrimaryKey=false, isUniqueKey=false, fieldPosition=3}, {fieldName='phonenumber', rawDataTypeNum=3, isPrimaryKey=false, isUniqueKey=false, fieldPosition=4}, {fieldName='address', rawDataTypeNum=15, isPrimaryKey=false, isUniqueKey=false, fieldPosition=5}], databaseName='dtstestdata', tableName='order', primaryIndexInfo [{indexType=PrimaryKey, indexFields=[{fieldName='orderid', rawDataTypeNum=3, isPrimaryKey=true, isUniqueKey=false, fieldPosition=0}], cardinality=0, nullable=true, isFirstUniqueIndex=false, name=null}], uniqueIndexInfo [[]], partitionFields = null}] Before image {[Field [orderid] [1] Field [username] [ Jane McLenahan] Field [ordertime] [xxx] Field [commodity] [xxx] Field [phonenumber] [ xxx] Field [address] [xxx] ]} After image {[Field [orderid] [1] Field [username] [ Jane] Field [ordertime] [xxx] Field [commodity] [xxx] Field [phonenumber] [ xxx] Field [address] [xxx] ]} (com.aliyun.dts.subscribe.clients.recordprocessor.DefaultRecordPrintListener) -
Klien SDK secara berkala mengumpulkan dan menampilkan statistik tentang konsumsi data, termasuk jumlah total dan volume catatan yang dikirim serta diterima, dan permintaan per detik (RPS).
[2021-05-18 16:25:09,167] INFO {"outCounts":488616.0,"outBytes":48606134,"outRps":1.15,"outBps":114.57,"count":11.0,"inBytes":60118961,"DStoreRecordQueue":0.0,"inCounts":557154.0,"inRps":1.12,"inBps":112.44,"__dt":1621326309167,"DefaultUserRecordQueue":0.0} (log_metrics)Tabel 2. Statistik konsumsi data
Parameter
Deskripsi
outCountsJumlah total catatan data yang dikonsumsi oleh klien SDK.
outBytesVolume total data yang dikonsumsi oleh klien SDK. Satuan: byte.
outRpsJumlah permintaan per detik (RPS) saat klien SDK mengonsumsi data.
outBpsLaju konsumsi data klien SDK, dalam bit per detik (bps).
inBytesVolume total data yang dikirim oleh server DTS. Satuan: byte.
DStoreRecordQueueUkuran antrian cache data internal untuk catatan masuk dari server DTS.
inCountsJumlah total catatan data yang dikirim oleh server DTS.
inRpsRPS saat server DTS mengirim data.
__dtStempel waktu saat klien SDK menerima data. Satuan: milidetik.
DefaultUserRecordQueueUkuran antrian data yang menyimpan catatan siap diproses oleh aplikasi konsumen.
-
Simpan dan kueri titik pemeriksaan konsumsi
Untuk memulai atau melanjutkan konsumsi data (misalnya, saat peluncuran pertama, restart, atau retry internal), klien SDK memerlukan titik pemeriksaan konsumsi. Tabel berikut menjelaskan cara mengelola dan mengkueri checkpoint dalam berbagai skenario untuk mencegah kehilangan data, meminimalkan konsumsi duplikat, dan memungkinkan konsumsi sesuai kebutuhan.
|
Skenario |
Mode penggunaan SDK |
Metode kueri |
|
Kueri titik pemeriksaan konsumsi |
Mode ASSIGN, mode SUBSCRIBE |
|
|
Startup awal: Memberikan checkpoint untuk memulai konsumsi. |
Mode ASSIGN, mode SUBSCRIBE |
Berdasarkan mode penggunaan klien SDK, pilih file DTSConsumerAssignDemo.java atau DTSConsumerSubscribeDemo.java, lalu konfigurasikan parameter |
|
Klien SDK perlu memberikan kembali titik pemeriksaan konsumsi terakhir yang dicatat untuk melanjutkan konsumsi setelah retry internal. |
Mode ASSIGN |
Cari titik pemeriksaan konsumsi terakhir yang dicatat dengan urutan berikut. Pencarian berhenti dan mengembalikan informasi checkpoint segera setelah ditemukan:
|
|
Mode SUBSCRIBE |
Cari titik pemeriksaan konsumsi terakhir yang dicatat dengan urutan berikut. Pencarian berhenti dan mengembalikan informasi checkpoint segera setelah ditemukan:
|
|
|
Klien SDK telah di-restart dan perlu memberikan kembali titik pemeriksaan konsumsi terakhir yang dicatat untuk melanjutkan konsumsi. |
Mode ASSIGN |
Kueri titik pemeriksaan konsumsi berdasarkan konfigurasi
|
|
Mode SUBSCRIBE |
Dalam mode ini, konfigurasi
|
Persist titik pemeriksaan konsumsi
Saat terjadi kejadian pemulihan bencana pada modul pengumpulan data inkremental (terutama dalam mode SUBSCRIBE), modul baru tidak menyimpan titik pemeriksaan konsumsi terbaru dari klien. Klien mungkin melanjutkan dari checkpoint yang lebih lama, sehingga menyebabkan konsumsi duplikat data historis. Misalnya, sebelum alih bencana, rentang checkpoint modul lama adalah dari 08:00:00 pada 11 November 2023 hingga 08:00:00 pada 12 November 2023, dan checkpoint klien adalah 08:00:00 pada 12 November 2023. Setelah alih bencana, rentang checkpoint modul baru adalah dari 10:00:00 pada 8 November 2023 hingga 08:01:00 pada 12 November 2023. Klien mulai dari checkpoint awal modul baru (10:00:00 pada 8 November 2023), sehingga menyebabkan konsumsi duplikat.
Untuk menghindari konsumsi duplikat dalam skenario ini, konfigurasikan penyimpanan checkpoint persisten di sisi klien. Contoh berikut menyediakan salah satu implementasi yang mungkin yang dapat Anda sesuaikan dengan kebutuhan Anda.
-
Buat kelas
UserMetaStoreyang memperluasAbstractUserMetaStore.Sebagai contoh, untuk menyimpan informasi checkpoint dalam database MySQL, gunakan kode Java berikut:
public class UserMetaStore extends AbstractUserMetaStore { @Override protected void saveData(String groupID, String toStoreJson) { Connection con = getConnection(); String sql = "insert into dts_checkpoint(group_id, checkpoint) values(?, ?)"; PreparedStatement pres = null; ResultSet rs = null; try { pres = con.prepareStatement(sql); pres.setString(1, groupID); pres.setString(2, toStoreJson); pres.execute(); } catch (Exception e) { e.printStackTrace(); } finally { close(rs, pres, con); } } @Override protected String getData(String groupID) { Connection con = getConnection(); String sql = "select checkpoint from dts_checkpoint where group_id = ?"; PreparedStatement pres = null; ResultSet rs = null; try { pres = con.prepareStatement(sql); pres.setString(1, groupID); ResultSet rs = pres.executeQuery() String checkpoint = rs.getString("checkpoint"); return checkpoint; } catch (Exception e) { e.printStackTrace(); } finally { close(rs, pres, con); } } } -
Dalam file consumerContext.java, konfigurasikan media penyimpanan eksternal menggunakan metode
setUserRegisteredStore(new UserMetaStore()).
FAQ
-
Bagaimana cara mengatasi masalah koneksi dengan instans Langganan Data?
Atasi masalah tersebut berdasarkan pesan error. Untuk informasi selengkapnya, lihat Troubleshooting.
-
Dalam format apa titik pemeriksaan konsumsi dipersist?
Data titik pemeriksaan konsumsi yang dipersist disimpan dalam format JSON. Titik pemeriksaan yang dipersist adalah stempel waktu UNIX yang dapat diberikan langsung ke SDK. Dalam tanggapan contoh berikut, nilai
1700709977untuk kunci"timestamp"adalah titik pemeriksaan konsumsi yang dipersist.{"groupID":"dtsglg11d48230***","streamCheckpoint":[{"partition":0,"offset":577989,"topic":"ap_southeast_1_vpc_rm_t4n22s21iysr6****_root_version2","timestamp":1700709977,"info":""}]}
Troubleshooting
|
Masalah |
Pesan error |
Penyebab |
Solusi |
|
Tidak dapat terhubung |
|
|
Masukkan nilai yang benar untuk parameter |
|
Alamat broker tidak dapat terhubung ke alamat IP aktual. |
||
|
Username atau kata sandi salah. |
||
|
Dalam file consumerContext.java, parameter |
Berikan titik pemeriksaan konsumsi yang berada dalam rentang data instans Langganan Data. Untuk informasi selengkapnya, lihat Parameter yang diperlukan. |
|
|
Konsumsi melambat |
N/A |
|
|