Topik ini menjelaskan cara menggunakan demo client Kafka untuk mengonsumsi data yang dilacak. Fitur pelacakan perubahan versi baru memungkinkan Anda mengonsumsi data yang dilacak dengan client Kafka versi V0.11 hingga V2.7.
Catatan penggunaan
-
Jika Anda mengaktifkan auto commit saat menggunakan fitur pelacakan perubahan, beberapa data mungkin dikomit sebelum dikonsumsi, sehingga menyebabkan kehilangan data. Kami menyarankan agar Anda melakukan commit data secara manual.
CatatanJika data gagal dikomit, Anda dapat me-restart client untuk melanjutkan konsumsi data dari titik pemeriksaan konsumsi terakhir yang terekam. Namun, data duplikat mungkin dihasilkan selama periode tersebut. Anda harus memfilter data duplikat tersebut secara manual.
-
Data diserialisasi dan disimpan dalam format Avro. Untuk informasi lebih lanjut, lihat Record.avsc.
PeringatanJika Anda tidak menggunakan client Kafka yang dijelaskan dalam topik ini, Anda harus mengurai data yang dilacak berdasarkan skema Avro (Contoh deserialisasi Avro DTS) dan memeriksa data yang telah diurai.
-
Satuan pencarian adalah detik saat Data Transmission Service (DTS) memanggil operasi
offsetForTimes, sedangkan satuan pencarian adalah milidetik saat klien Kafka bawaan memanggil operasi tersebut. -
Koneksi sementara (transient connections) mungkin terjadi antara client Kafka dan server pelacakan perubahan karena berbagai alasan, seperti disaster recovery. Jika Anda tidak menggunakan client Kafka yang dijelaskan dalam topik ini, client Kafka Anda harus memiliki kemampuan koneksi ulang jaringan.
-
Jika Anda menggunakan client Kafka native untuk mengonsumsi data yang dilacak, modul pengumpulan data inkremental mungkin berubah di DTS. Dalam mode subscribe, titik pemeriksaan konsumsi yang disimpan oleh client Kafka ke server DTS akan dihapus. Anda perlu menentukan titik pemeriksaan konsumsi sesuai kebutuhan bisnis Anda untuk mengonsumsi data yang dilacak. Jika ingin mengonsumsi data dalam mode subscribe, kami menyarankan agar Anda menggunakan demo SDK yang disediakan oleh DTS untuk melacak dan mengonsumsi data, atau mengelola titik pemeriksaan konsumsi secara manual. Untuk informasi lebih lanjut, lihat Mengonsumsi data berlangganan menggunakan SDK dan bagian Mengelola titik pemeriksaan konsumsi pada topik ini.
Jalankan client Kafka
Unduh demo client Kafka. Untuk informasi lebih lanjut tentang cara menggunakan demo ini, lihat Readme.
-
Klik
dan pilih Download ZIP untuk mengunduh paketnya. -
Jika Anda menggunakan client Kafka versi 2.0, Anda harus mengubah nomor versi dalam file subscribe_example-master/javaimpl/pom.xml menjadi 2.0.0.
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.0.0</version>
</dependency>
|
Langkah |
Direktori atau file terkait |
|
1. Gunakan consumer Kafka native untuk mendapatkan data inkremental dari instans pelacakan perubahan. |
subscribe_example-master/javaimpl/src/main/java/recordgenerator/ |
|
2. Deserialisasi gambar data inkremental, lalu peroleh pre-image , post-image , dan atribut lainnya. Peringatan
|
subscribe_example-master/javaimpl/src/main/java/boot/RecordPrinter.java |
|
3. Konversi nilai dataTypeNumber dalam data yang telah dideserialisasi menjadi tipe data database yang sesuai. Catatan
Untuk informasi lebih lanjut, lihat bagian berikut dalam topik ini: |
subscribe_example-master/javaimpl/src/main/java/recordprocessor/mysql/ |
Prosedur
Langkah-langkah berikut menunjukkan cara menjalankan client Kafka untuk mengonsumsi data yang dilacak. Dalam contoh ini, digunakan IntelliJ IDEA Community Edition 2018.1.4 untuk Windows.
-
Buat instans pelacakan perubahan. Untuk informasi lebih lanjut, lihat Ikhtisar pelacakan perubahan.
-
Buat satu atau beberapa kelompok konsumen. Untuk informasi lebih lanjut, lihat Menambahkan kelompok konsumen.
-
Unduh paket demo client Kafka dan ekstrak paket tersebut.
CatatanKlik
dan pilih Download ZIP untuk mengunduh paketnya. -
Buka IntelliJ IDEA. Di jendela yang muncul, klik Open.
-
Di kotak dialog yang muncul, buka direktori tempat demo yang diunduh berada. Temukan file pom.xml.
Buka
kafkademo>subscribe_example-master>javaimpl, pilihpom.xml, lalu klik OK. -
Di kotak dialog yang muncul, pilih Open as Project.
-
Di jendela tool Project IntelliJ IDEA, klik folder-folder untuk menemukan file demo client Kafka, lalu klik ganda file tersebut. Nama file-nya adalah NotifyDemoDB.java.
-
Tentukan parameter dalam file NotifyDemoDB.java.
public static Properties getConfigs() { Properties properties = new Properties(); // user password and sid for auth properties.setProperty(USER_NAME, "dtstest"); properties.setProperty(PASSWORD_NAME, "xxx"); properties.setProperty(SID_NAME, "dtsxxx"); // kafka consumer group general same with sid properties.setProperty(GROUP_NAME, "dtsxxx"); // topic to consume, partition is 0 properties.setProperty(KAFKA_TOPIC, "cn_hangzhou_xxx"); // kafka broker url properties.setProperty(KAFKA_BROKER_URL_NAME, "dts-cn-xxx.com:18001"); // initial checkpoint for first seek(a timestamp to set, eg 1566180200 if you want (Mon Aug 19 10:03:21 CST 2019)) properties.setProperty(INITIAL_CHECKPOINT_NAME, "1583307907"); // if force use config checkpoint when start. for checkpoint reset properties.setProperty(USE_CONFIG_CHECKPOINT_NAME, "true"); // use consumer assign or subscribe interface // when use subscribe mode, group config is required. kafka consumer group is enabled properties.setProperty(SUBSCRIBE_MODE_NAME, "assign"); return properties; }Parameter
Deskripsi
Cara memperoleh nilai parameter
USER_NAME
Username akun kelompok konsumen.
PeringatanJika Anda tidak menggunakan client Kafka yang dijelaskan dalam topik ini, Anda harus menentukan parameter ini dalam format berikut:
<Username>-<Consumer group ID>. Contoh:dtstest-dtsae******bpv. Jika tidak, koneksi akan gagal.Di Konsol DTS, temukan instans pelacakan perubahan yang ingin Anda kelola lalu klik ID instans tersebut. Di panel navigasi sebelah kiri, klik Consume Data. Di halaman yang muncul, Anda dapat melihat informasi tentang kelompok konsumen, seperti ID atau nama dan akun kelompok konsumen tersebut.
CatatanPassword akun kelompok konsumen ditentukan saat Anda membuat kelompok konsumen tersebut.
PASSWORD_NAME
Password akun tersebut.
SID_NAME
ID kelompok konsumen.
GROUP_NAME
Nama kelompok konsumen. Atur parameter ini ke ID kelompok konsumen.
KAFKA_TOPIC
Nama topik yang dilacak dari instans pelacakan perubahan.
Di Konsol DTS, temukan instans pelacakan perubahan yang ingin Anda kelola lalu klik ID instans tersebut. Di halaman Basic Information, Anda dapat melihat informasi tentang topik dan jaringan. Di bagian Basic Information pada halaman detail tugas pelacakan perubahan, peroleh nilai Subscription Topic. Di bagian Network, peroleh VPC network address (contoh format:
xxx.aliyuncs.com:18003).KAFKA_BROKER_URL_NAME
Titik akhir instans pelacakan perubahan.
CatatanJika Anda melacak perubahan data melalui jaringan internal, latensi jaringan minimal. Ini berlaku jika instans Elastic Compute Service (ECS) tempat Anda men-deploy client Kafka berada di jaringan klasik atau dalam virtual private cloud (VPC) yang sama dengan instans pelacakan perubahan.
INITIAL_CHECKPOINT_NAME
Titik pemeriksaan konsumsi data yang telah dikonsumsi. Nilainya berupa stempel waktu UNIX. Contoh: 1592269238.
Catatan-
Anda harus menyimpan titik pemeriksaan konsumsi karena alasan berikut:
-
Jika proses konsumsi terganggu, Anda dapat menentukan titik pemeriksaan konsumsi pada client Kafka untuk melanjutkan konsumsi data. Hal ini mencegah kehilangan data.
-
Saat memulai client Kafka, Anda dapat menentukan titik pemeriksaan konsumsi untuk mengonsumsi data sesuai kebutuhan bisnis Anda.
-
-
Jika parameter SUBSCRIBE_MODE_NAME diatur ke subscribe, parameter INITIAL_CHECKPOINT_NAME yang Anda tentukan hanya berlaku saat pertama kali menjalankan client Kafka.
Titik pemeriksaan konsumsi data yang telah dikonsumsi harus berada dalam rentang data instans pelacakan perubahan. Titik pemeriksaan konsumsi harus dikonversi menjadi stempel waktu UNIX. Di daftar tugas pelacakan perubahan DTS, Anda dapat melihat bidang Data Range dari tugas yang sesuai, yang menunjukkan waktu mulai dan akhir data yang dapat dikonsumsi. Anda dapat menggunakan bidang ini untuk menentukan rentang nilai valid INITIAL_CHECKPOINT_NAME.
Catatan-
Anda dapat melihat rentang data instans pelacakan perubahan di kolom Data Range pada halaman Tugas Pelacakan Perubahan.
-
Anda dapat menggunakan mesin pencari untuk memperoleh konverter stempel waktu UNIX.
USE_CONFIG_CHECKPOINT_NAME
Menentukan apakah client dipaksa mengonsumsi data dari titik pemeriksaan konsumsi yang ditentukan. Nilai default: true. Anda dapat mengatur parameter ini ke true untuk mencegah data yang telah diterima tetapi belum diproses hilang.
Tidak ada
SUBSCRIBE_MODE_NAME
Menentukan apakah akan menjalankan dua atau lebih client Kafka untuk satu kelompok konsumen. Jika ingin menggunakan fitur ini, atur parameter ini ke subscribe untuk client Kafka tersebut.
Nilai default-nya adalah assign, yang menunjukkan bahwa fitur ini tidak digunakan. Kami menyarankan agar Anda hanya men-deploy satu client Kafka untuk satu kelompok konsumen.
Tidak ada
-
-
Di bilah menu atas IntelliJ IDEA, pilih untuk menjalankan client.
CatatanJika Anda menjalankan IntelliJ IDEA untuk pertama kalinya, diperlukan waktu tertentu untuk memuat dan menginstal dependensi terkait.
Hasil pada client Kafka
Hasil berikut menunjukkan bahwa client Kafka dapat melacak perubahan data dari database sumber.
[2020-03-09 10:41:52,408] INFO [Consumer clientId=consumer-1, groupId=dts_xxx] Discovered coordinator xxx (id: xxx rack: null) (org.apache.kafka.clients.consumer.internals.AbstractCoordinator)
[2020-03-09 10:41:57,203] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721711, offset: 1732521, info: 1583721711] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:41:57,571] INFO EtlRecordProcessor: haven't receive records from generator for 5s (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:02,203] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721721, offset: 1732539, info: 1583721721] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:07,204] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721726, offset: 1732544, info: 1583721726] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:12,205] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721731, offset: 1732548, info: 1583721731] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:17,205] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721736, offset: 1732554, info: 1583721736] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:22,205] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721741, offset: 1732559, info: 1583721741] (recordprocessor.EtlRecordProcessor)
[2020-03-09 10:42:27,206] INFO commit record with checkpoint Checkpoint[ topicPartition: cn_hangzhou_rm_xxx_dtstest-0timestamp: 1583721746, offset: 1732569, info: 1583721746] (recordprocessor.EtlRecordProcessor)
Anda dapat menghapus garis miring ganda (//) dari string //log.info(ret) di baris 25 file NotifyDemoDB.java. Kemudian, jalankan kembali client untuk melihat informasi perubahan data.
FAQ
-
T: Mengapa saya perlu mencatat titik pemeriksaan konsumsi client Kafka?
J: Titik pemeriksaan konsumsi yang dicatat oleh DTS adalah titik waktu saat DTS menerima operasi commit dari client Kafka. Titik pemeriksaan konsumsi yang terekam mungkin berbeda dari waktu konsumsi aktual. Jika aplikasi bisnis atau client Kafka terputus secara tak terduga, Anda dapat menentukan titik pemeriksaan konsumsi yang akurat untuk melanjutkan konsumsi data. Hal ini mencegah kehilangan data atau konsumsi data duplikat.
Mengelola titik pemeriksaan konsumsi
-
Konfigurasikan client Kafka untuk mendengarkan alih bencana modul pengumpulan data di DTS.
Anda dapat mengonfigurasi properti consumer client Kafka untuk mendengarkan alih bencana modul pengumpulan data di DTS. Kode berikut memberikan contoh cara mengonfigurasi properti consumer:
properties.setProperty(ConsumerConfig.INTERCEPTOR_CLASSES_CONFIG, ClusterSwitchListener.class.getName());Kode berikut memberikan contoh cara mengimplementasikan ClusterSwitchListener:
public class ClusterSwitchListener implements ClusterResourceListener, ConsumerInterceptor { private final static Logger LOG = LoggerFactory.getLogger(ClusterSwitchListener.class); private ClusterResource originClusterResource = null; private ClusterResource currentClusterResource = null; public ConsumerRecords onConsume(ConsumerRecords records) { return records; } public void close() { } public void onCommit(Map offsets) { } public void onUpdate(ClusterResource clusterResource) { synchronized (this) { originClusterResource = currentClusterResource; currentClusterResource = clusterResource; if (null == originClusterResource) { LOG.info("Cluster updated to " + currentClusterResource.clusterId()); } else { if (originClusterResource.clusterId().equals(currentClusterResource.clusterId())) { LOG.info("Cluster not changed on update:" + clusterResource.clusterId()); } else { LOG.error("Cluster changed"); throw new ClusterSwitchException("Cluster changed from " + originClusterResource.clusterId() + " to " + currentClusterResource.clusterId() + ", consumer require restart"); } } } } public boolean isClusterResourceChanged() { if (null == originClusterResource) { return false; } if (originClusterResource.clusterId().equals(currentClusterResource.clusterId())) { return false; } return true; } public void configure(Map<String, ?> configs) { } public static class ClusterSwitchException extends KafkaException { public ClusterSwitchException(String message, Throwable cause) { super(message, cause); } public ClusterSwitchException(String message) { super(message); } public ClusterSwitchException(Throwable cause) { super(cause); } public ClusterSwitchException() { super(); } } -
Tentukan titik pemeriksaan konsumsi berdasarkan alih bencana modul pengumpulan data di DTS yang terdeteksi.
Atur titik pemeriksaan konsumsi awal pelacakan data berikutnya ke stempel waktu entri data terlacak terbaru yang dikonsumsi oleh client. Kode berikut memberikan contoh cara menentukan titik pemeriksaan konsumsi:
try{ //do some action } catch (ClusterSwitchListener.ClusterSwitchException e) { reset(); } // Reset the consumption checkpoint. public reset() { long offset = kafkaConsumer.offsetsForTimes(timestamp); kafkaConsumer.seek(tp,offset); }CatatanUntuk informasi lebih lanjut tentang contoh-contoh tersebut, lihat KafkaRecordFetcher.
Pemetaan antara tipe data MySQL dan nilai dataTypeNumber
Untuk informasi lebih lanjut, lihat SQL Type field.
|
Tipe data MySQL |
Nilai dataTypeNumber |
|
MYSQL_TYPE_DECIMAL |
0 |
|
MYSQL_TYPE_INT8 |
1 |
|
MYSQL_TYPE_INT16 |
2 |
|
MYSQL_TYPE_INT32 |
3 |
|
MYSQL_TYPE_FLOAT |
4 |
|
MYSQL_TYPE_DOUBLE |
5 |
|
MYSQL_TYPE_NULL |
6 |
|
MYSQL_TYPE_TIMESTAMP |
7 |
|
MYSQL_TYPE_INT64 |
8 |
|
MYSQL_TYPE_INT24 |
9 |
|
MYSQL_TYPE_DATE |
10 |
|
MYSQL_TYPE_TIME |
11 |
|
MYSQL_TYPE_DATETIME |
12 |
|
MYSQL_TYPE_YEAR |
13 |
|
MYSQL_TYPE_DATE_NEW |
14 |
|
MYSQL_TYPE_VARCHAR |
15 |
|
MYSQL_TYPE_BIT |
16 |
|
MYSQL_TYPE_TIMESTAMP_NEW |
17 |
|
MYSQL_TYPE_DATETIME_NEW |
18 |
|
MYSQL_TYPE_TIME_NEW |
19 |
|
MYSQL_TYPE_JSON |
245 |
|
MYSQL_TYPE_DECIMAL_NEW |
246 |
|
MYSQL_TYPE_ENUM |
247 |
|
MYSQL_TYPE_SET |
248 |
|
MYSQL_TYPE_TINY_BLOB |
249 |
|
MYSQL_TYPE_MEDIUM_BLOB |
250 |
|
MYSQL_TYPE_LONG_BLOB |
251 |
|
MYSQL_TYPE_BLOB |
252 |
|
MYSQL_TYPE_VAR_STRING |
253 |
|
MYSQL_TYPE_STRING |
254 |
|
MYSQL_TYPE_GEOMETRY |
255 |
Pemetaan antara tipe data Oracle dan nilai dataTypeNumber
|
Tipe data Oracle |
Nilai dataTypeNumber |
|
VARCHAR2/NVARCHAR2 |
1 |
|
NUMBER/FLOAT |
2 |
|
LONG |
8 |
|
DATE |
12 |
|
RAW |
23 |
|
LONG_RAW |
24 |
|
UNDEFINED |
29 |
|
XMLTYPE |
58 |
|
ROWID |
69 |
|
CHAR and NCHAR |
96 |
|
BINARY_FLOAT |
100 |
|
BINARY_DOUBLE |
101 |
|
CLOB/NCLOB |
112 |
|
BLOB |
113 |
|
BFILE |
114 |
|
TIMESTAMP |
180 |
|
TIMESTAMP_WITH_TIME_ZONE |
181 |
|
INTERVAL_YEAR_TO_MONTH |
182 |
|
INTERVAL_DAY_TO_SECOND |
183 |
|
UROWID |
208 |
|
TIMESTAMP_WITH_LOCAL_TIME_ZONE |
231 |
Pemetaan antara tipe data PostgreSQL dan nilai dataTypeNumber
|
Tipe data PostgreSQL |
Nilai dataTypeNumber |
|
INT2/SMALLINT |
21 |
|
INT4/INTEGER/SERIAL |
23 |
|
INT8/BIGINT |
20 |
|
CHARACTER |
18 |
|
CHARACTER VARYING |
1043 |
|
REAL |
700 |
|
DOUBLE PRECISION |
701 |
|
NUMERIC |
1700 |
|
MONEY |
790 |
|
DATE |
1082 |
|
TIME/TIME WITHOUT TIME ZONE |
1083 |
|
TIME WITH TIME ZONE |
1266 |
|
TIMESTAMP/TIMESTAMP WITHOUT TIME ZONE |
1114 |
|
TIMESTAMP WITH TIME ZONE |
1184 |
|
BYTEA |
17 |
|
TEXT |
25 |
|
JSON |
114 |
|
JSONB |
3082 |
|
XML |
142 |
|
UUID |
2950 |
|
POINT |
600 |
|
LSEG |
601 |
|
PATH |
602 |
|
BOX |
603 |
|
POLYGON |
604 |
|
LINE |
628 |
|
CIDR |
650 |
|
CIRCLE |
718 |
|
MACADDR |
829 |
|
INET |
869 |
|
INTERVAL |
1186 |
|
TXID_SNAPSHOT |
2970 |
|
PG_LSN |
3220 |
|
TSVECTOR |
3614 |
|
TSQUERY |
3615 |