All Products
Search
Document Center

Data Transmission Service:Gunakan client Kafka untuk mengonsumsi data yang dilacak

Last Updated:Aug 25, 2026

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.

    Catatan

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

    Peringatan

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

Catatan
  • Klik code 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>
Tabel 1 Deskripsi proses

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
  • Jika instans sumber adalah database Oracle yang dikelola sendiri, Anda harus mengaktifkan supplemental logging untuk semua kolom. Hal ini memastikan client dapat berhasil mengonsumsi data yang dilacak serta menjamin integritas pre-image dan post-image.

  • Jika instans sumber bukan database Oracle yang dikelola sendiri, DTS tidak menjamin integritas pre-image. Kami menyarankan agar Anda memverifikasi pre-image yang diperoleh.

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.

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.

  1. Buat instans pelacakan perubahan. Untuk informasi lebih lanjut, lihat Ikhtisar pelacakan perubahan.

  2. Buat satu atau beberapa kelompok konsumen. Untuk informasi lebih lanjut, lihat Menambahkan kelompok konsumen.

  3. Unduh paket demo client Kafka dan ekstrak paket tersebut.

    Catatan

    Klik code dan pilih Download ZIP untuk mengunduh paketnya.

  4. Buka IntelliJ IDEA. Di jendela yang muncul, klik Open.

  5. Di kotak dialog yang muncul, buka direktori tempat demo yang diunduh berada. Temukan file pom.xml.

    Buka kafkademo > subscribe_example-master > javaimpl, pilih pom.xml, lalu klik OK.

  6. Di kotak dialog yang muncul, pilih Open as Project.

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

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

    Peringatan

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

    Catatan

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

    Catatan

    Jika 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

  9. Di bilah menu atas IntelliJ IDEA, pilih Run > Run untuk menjalankan client.

    Catatan

    Jika 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

  1. 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();
            }
    
        }
  2. 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);
    }
    Catatan

    Untuk 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