Buat saluran langganan data dalam mode pull dan push.
Prasyarat
Alamat IP client Lindorm dan client Message Queue for Apache Kafka Anda telah ditambahkan ke daftar putih instans Lindorm Anda. Untuk informasi selengkapnya, lihat Konfigurasi daftar putih.
Instans Lindorm sumber dan instans Message Queue for Apache Kafka tujuan telah terhubung ke Lindorm Tunnel Service (LTS). Untuk informasi selengkapnya, lihat Membangun koneksi jaringan.
Sumber data LindormTable telah dibuat. Untuk informasi selengkapnya, lihat Tambahkan sumber data LindormTable.
Sumber data Kafka telah dibuat. Untuk informasi lebih lanjut, lihat Tambahkan sumber data Kafka.
Fitur langganan data atau pelacakan perubahan telah diaktifkan. Untuk informasi selengkapnya, lihat Aktifkan langganan data atau Aktifkan pelacakan perubahan.
Buat saluran langganan data mode pull
Prosedur
-
Buka halaman LTS. Di panel navigasi sebelah kiri, pilih Change Tracking > Pull mode.
Anda akan diarahkan ke halaman Subscription channel list dari Lindorm CDC. Tombol Create data subscription channel berada di pojok kanan atas. Daftar tersebut menampilkan ID saluran langganan, nama tabel Lindorm, nama topik, dan aksi (Details dan Delete) untuk saluran yang sudah ada.
-
Klik Create data subscription channel dan konfigurasikan parameter berikut.
Parameter
Deskripsi
Source cluster
Masukkan ID instans Lindorm.
Lindorm table name
Pilih tabel Lindorm yang ingin Anda langganan. Setiap saluran hanya dapat berlangganan pada satu tabel.
Topic name
Nama topik konsumsi data.
Data expiration time (days)
Jumlah hari penyimpanan data. Nilai default adalah 7.
Number of topic partitions
Jumlah partisi untuk topik. Beberapa partisi memungkinkan konsumsi data secara konkuren. Nilai default adalah 4.
-
Klik Submit.
-
(Opsional) Untuk melihat detail saluran, temukan saluran tersebut dalam daftar lalu klik Details di kolom Operation. Anda dapat melihat detail saluran, detail konsumsi, dan detail penyimpanan.
-
(Opsional) Gunakan contoh kode berikut dengan client Kafka untuk mengonsumsi data yang telah dilanggan.
import org.apache.hadoop.hbase.util.Bytes; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import java.time.Duration; import java.util.Arrays; import java.util.Properties; public class TestConsume { public static void main(String[] args) throws Exception { // Nama topik yang Anda tentukan saat membuat saluran langganan data. String topic = "test-topic"; // Properti untuk menghubungkan ke titik akhir. Properties props = new Properties(); // Tentukan alamat titik akhir. props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "ld-xxx:9092"); // Deserializer kunci. Jangan ubah nilai ini. props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName()); // Deserializer nilai. Jangan ubah nilai ini. props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName()); // Nama kelompok konsumen. Kelompok ini dibuat otomatis saat konsumsi. props.put(ConsumerConfig.GROUP_ID_CONFIG, "group-id-0"); // Buat konsumen. KafkaConsumer<byte[], byte[]> consumer = new KafkaConsumer<>(props); // Berlangganan ke topik. consumer.subscribe(Arrays.asList(topic)); // Tarik data menggunakan konsumen. ConsumerRecords<byte[], byte[]> records = consumer.poll(Duration.ofMillis(10000)); for (ConsumerRecord<byte[], byte[]> record : records) { // Lihat isi data. System.out.println("key: " + Bytes.toString(record.key())); System.out.println("value: " + Bytes.toString(record.value())); } // Commit offset konsumen saat ini. consumer.commitSync(); // Tutup konsumen. consumer.close(); } }CatatanUntuk informasi selengkapnya mengenai format konsumsi data, lihat Format konsumsi data.
Buat saluran Push untuk langganan data
Proses
Gambar berikut menunjukkan bagaimana paket data inkremental dari tabel Lindorm didorong ke Message Queue for Apache Kafka.
Buat aliran Lindorm
-
Buka Konsol Lindorm Tunnel Service (LTS). Di panel navigasi sebelah kiri, pilih Change Tracking > Push.
-
Klik create dan konfigurasikan parameter.
Parameter
Deskripsi
Lindorm Cluster
Pilih sumber data LindormTable yang telah dibuat.
Table Name
Nama tabel tempat Anda ingin menangkap perubahan data. Formatnya adalah
namespace.tablename.Contohnya,
ns1.table1berarti Anda ingin menangkap data dari tabeltable1di namespacens1.Blacklist table (Optional)
Tentukan tabel yang datanya tidak boleh didorong. Perubahan pada tabel-tabel ini akan diabaikan.
MessageStorage Type
Pilih KAFKA.
Storage Datasource
Pilih sumber data Kafka yang telah dibuat.
MessageStorage Config
-
kafka_topic: Tentukan nama topik Kafka. -
kafka_ttl: Biarkan kosong. -
kafka_partition_num: Biarkan kosong.
PentingDalam mode push, Anda harus membuat topik di Kafka terlebih dahulu.
MessageVersion
Format pesan. Default-nya adalah DebeziumV2.
Message Config
-
old_image: Menentukan apakah pesan mencakup nilai baris sebelum perubahan. Anda harus mengatur parameter ini ketrue. Mengaturnya kefalsetidak didukung. -
new_image: Menentukan apakah pesan mencakup nilai baris setelah perubahan. Anda harus mengatur parameter ini ketrue. Mengaturnya kefalsetidak didukung. -
with_schema: Menentukan apakah pesan mencakup skema tabel. Atur parameter ini kefalseuntuk mencegah ukuran pesan terlalu besar. -
ignore_family_prefix: Menentukan apakah awalan keluarga kolom dihapus dari nama kolom yang diekspor. Misalnya, jika nama kolom lengkapnya adalahf:name, mengatur parameter ini ketrueakan menghasilkan nama kolom yang diekspor menjadiname.
-
-
Klik Submit.