Anda dapat menyinkronkan data secara real time dari ApsaraMQ for Kafka ke ApsaraDB for ClickHouse dengan menggunakan Mesin tabel Kafka bawaan dan Tampilan yang di-materialisasi.
Batasan
Anda hanya dapat menyinkronkan data dari instans ApsaraMQ for Kafka dan kluster Kafka yang dikelola sendiri yang di-deploy pada Instance ECS.
Prasyarat
-
ApsaraDB for ClickHouse:
-
Kluster tujuan telah dibuat di wilayah dan VPC yang sama dengan instans ApsaraMQ for Kafka. Untuk informasi selengkapnya, lihat Buat kluster.
-
Akun database dengan izin yang diperlukan telah dibuat untuk kluster tujuan. Untuk informasi selengkapnya, lihat Manajemen Akun.
-
-
ApsaraMQ for Kafka:
Catatan penggunaan
-
Topik yang berlangganan oleh tabel eksternal Kafka ApsaraDB for ClickHouse tidak boleh memiliki konsumen lain.
-
Saat membuat tabel eksternal Kafka, Tampilan yang di-materialisasi, dan tabel lokal, tipe bidang dari ketiga tabel tersebut harus sesuai.
Prosedur
Contoh berikut menyinkronkan data dari ApsaraMQ for Kafka ke tabel terdistribusi kafka_table_distributed di database default kluster Edisi Kompatibel Komunitas ApsaraDB for ClickHouse.
Langkah 1: Pahami cara kerja sinkronisasi
ApsaraDB for ClickHouse menggunakan Mesin tabel Kafka dan Tampilan yang di-materialisasi untuk mengonsumsi dan menyimpan data dari Kafka secara real time. Alur data adalah sebagai berikut.
-
Topik Kafka: Data sumber yang akan disinkronkan.
-
Tabel eksternal Kafka ApsaraDB for ClickHouse (tabel yang menggunakan Mesin tabel Kafka): menarik data sumber dari topik Kafka tertentu.
-
Tampilan yang di-materialisasi: membaca data sumber dari tabel eksternal Kafka dan memasukkan data tersebut ke tabel lokal di ApsaraDB for ClickHouse.
-
Tabel lokal: menyimpan data yang telah disinkronkan.
Langkah 2: Hubungkan ke kluster ApsaraDB for ClickHouse
Untuk informasi selengkapnya, lihat Hubungkan ke kluster ApsaraDB for ClickHouse menggunakan DMS.
Langkah 3: Buat tabel eksternal Kafka
Tabel eksternal Kafka menggunakan Mesin tabel Kafka untuk menarik data dari topik Kafka tertentu. Tabel ini memiliki karakteristik sebagai berikut:
-
Secara default, Anda tidak dapat langsung melakukan kueri terhadap tabel eksternal Kafka.
-
Tabel eksternal Kafka hanya digunakan untuk mengonsumsi data Kafka dan tidak menyimpan data. Anda harus menggunakan Tampilan yang di-materialisasi untuk memproses dan memasukkan data ke tabel tujuan.
Sintaks untuk membuat tabel adalah sebagai berikut.
Tipe bidang tabel eksternal Kafka harus konsisten dengan tipe data pesan di Kafka.
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
...
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'host:port1,host:port2,host:port3',
kafka_topic_list = 'topic_name1,topic_name2,...',
kafka_group_name = 'group_name',
kafka_format = 'data_format'[,]
[kafka_row_delimiter = 'delimiter_symbol',]
[kafka_num_consumers = N,]
[kafka_thread_per_consumer = 1,]
[kafka_max_block_size = 0,]
[kafka_skip_broken_messages = N,]
[kafka_commit_every_batch = 0,]
[kafka_auto_offset_reset = N]
Tabel berikut menjelaskan parameter umum.
|
Parameter |
Wajib |
Deskripsi |
|
kafka_broker_list |
Ya |
Daftar titik akhir broker untuk kluster Kafka yang dipisahkan koma. Untuk informasi selengkapnya tentang cara melihat titik akhir, lihat Lihat Titik Akhir.
|
|
kafka_topic_list |
Ya |
Daftar nama topik yang dipisahkan koma. Untuk informasi selengkapnya tentang cara melihat nama topik, lihat Buat topik. |
|
kafka_group_name |
Ya |
Nama kelompok konsumen Kafka. Untuk informasi selengkapnya, lihat Buat kelompok. |
|
kafka_format |
Ya |
Format isi pesan yang dapat diproses oleh ApsaraDB for ClickHouse. Catatan
Untuk informasi selengkapnya tentang format isi pesan yang didukung oleh ApsaraDB for ClickHouse, lihat Format untuk Data Input dan Output. |
|
kafka_row_delimiter |
Tidak |
Pemisah yang digunakan untuk memisahkan baris. Nilai default adalah \n. Anda juga dapat mengatur parameter ini agar sesuai dengan pemisah aktual yang digunakan dalam data Anda. |
|
kafka_num_consumers |
Tidak |
Jumlah konsumen untuk satu tabel. Nilai default adalah 1. Catatan
|
|
kafka_thread_per_consumer |
Tidak |
Menentukan apakah akan mengaktifkan thread khusus untuk setiap konsumen. Nilai default adalah 0. Nilai yang valid:
Untuk informasi selengkapnya tentang cara meningkatkan kecepatan konsumsi, lihat Penyetelan performa Kafka. |
|
kafka_max_block_size |
Tidak |
Ukuran maksimum, dalam byte, dari batch pesan Kafka. Nilai default adalah 65536. |
|
kafka_skip_broken_messages |
Tidak |
Jumlah error parsing yang diabaikan. Nilai default adalah 0. Jika Anda mengatur |
|
kafka_commit_every_batch |
Tidak |
Frekuensi commit Kafka. Nilai default adalah 0. Nilai yang valid:
|
|
kafka_auto_offset_reset |
Tidak |
Offset tempat pembacaan data Kafka dimulai. Nilai yang valid:
Catatan
Parameter ini tidak didukung untuk kluster ApsaraDB for ClickHouse yang menjalankan versi kernel 21.8. |
Untuk informasi selengkapnya tentang parameter, lihat Kafka.
Kode berikut memberikan contoh:
CREATE TABLE default.kafka_src_table ON CLUSTER `default`
(
-- Definisikan bidang skema tabel.
id Int32,
name String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092,alikafka-post-cn-****1-2-vpc.alikafka.aliyuncs.com:9092,alikafka-post-cn-****-3-vpc.alikafka.aliyuncs.com:9092',
kafka_topic_list = 'testforCK',
kafka_group_name = 'GroupForTestCK',
kafka_format = 'CSV';
Langkah 4: Buat tabel tujuan
Pilih pernyataan pembuatan tabel yang sesuai dengan edisi kluster Anda.
Untuk kluster Edisi Perusahaan, Anda hanya perlu membuat tabel lokal. Untuk kluster Edisi Kompatibel Komunitas, Anda mungkin perlu membuat tabel terdistribusi berdasarkan lingkungan dan kebutuhan Anda. Kode berikut memberikan contoh pernyataan. Untuk informasi selengkapnya tentang sintaks pembuatan tabel, lihat CREATE TABLE.
Edisi Perusahaan
CREATE TABLE default.kafka_table_local ON CLUSTER default (
id Int32,
name String
) ENGINE = MergeTree()
ORDER BY (id);
Jika Anda menerima error ON CLUSTER is not allowed for Replicated database saat menjalankan pernyataan ini, Anda dapat melakukan peningkatan versi kernel untuk mengatasi masalah tersebut. Untuk informasi selengkapnya tentang cara meningkatkan versi kernel, lihat Peningkatan versi mesin minor.
Edisi Kompatibel Komunitas
Mesin tabel untuk kluster single-replica dan double-replica berbeda. Pilih mesin yang sesuai berdasarkan tipe replica kluster Anda.
Saat membuat tabel di kluster dual-replica, Anda harus menggunakan mesin Replicated dari family mesin MergeTree. Jika Anda membuat tabel dengan mesin non-Replicated di kluster dual-replica, data tidak dapat direplikasi antar-replica, yang dapat menyebabkan inkonsistensi data.
Single-replica
-
Buat tabel lokal.
CREATE TABLE default.kafka_table_local ON CLUSTER default ( id Int32, name String ) ENGINE = MergeTree() ORDER BY (id); -
(Opsional) Buat tabel terdistribusi.
Jika Anda hanya perlu mengimpor data ke tabel lokal, lewati langkah ini.
Jika Anda memiliki kluster multi-node, kami merekomendasikan untuk membuat tabel terdistribusi.
CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local ENGINE = Distributed(default, default, kafka_table_local, id);
Double-replica
-
Buat tabel lokal.
CREATE TABLE default.kafka_table_local ON CLUSTER default ( id Int32, name String ) ENGINE = ReplicatedMergeTree() ORDER BY (id); -
(Opsional) Buat tabel terdistribusi.
Jika Anda hanya perlu mengimpor data ke tabel lokal, lewati langkah ini.
Jika Anda memiliki kluster multi-node, kami merekomendasikan untuk membuat tabel terdistribusi.
CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local ENGINE = Distributed(default, default, kafka_table_local, id);
Langkah 5: Buat Tampilan yang di-materialisasi
ApsaraDB for ClickHouse mengandalkan Tampilan yang di-materialisasi untuk membaca data sumber dari tabel eksternal Kafka dan memasukkan data tersebut ke dalam tabel lokal di ApsaraDB for ClickHouse.
Sintaks untuk membuat Tampilan yang di-materialisasi adalah sebagai berikut.
Pastikan bidang SELECT konsisten dengan struktur tabel tujuan, atau gunakan fungsi konversi agar format data sesuai dengan struktur tabel tujuan.
CREATE MATERIALIZED VIEW <view_name> ON CLUSTER default TO <dest_table> AS SELECT * FROM <src_table>;
Tabel berikut menjelaskan parameter-parameter tersebut.
|
Parameter |
Wajib |
Deskripsi |
Contoh |
|
view_name |
Ya |
Nama Tampilan. |
consumer |
|
dest_table |
Ya |
Tabel tujuan untuk menyimpan data Kafka.
|
|
|
src_table |
Ya |
Tabel eksternal Kafka. |
kafka_src_table |
Kode berikut memberikan contoh pernyataan.
Edisi Perusahaan
CREATE MATERIALIZED VIEW consumer ON CLUSTER default TO kafka_table_local AS SELECT * FROM kafka_src_table;
Edisi Kompatibel Komunitas
Dalam contoh ini, data sumber disimpan di tabel terdistribusi kafka_table_distributed.
CREATE MATERIALIZED VIEW consumer ON CLUSTER default TO kafka_table_distributed AS SELECT * FROM kafka_src_table;
Langkah 6: Verifikasi sinkronisasi
-
Kirim pesan ke topik dalam instans ApsaraMQ for Kafka.
-
Login ke Konsol ApsaraMQ for Kafka.
-
Pada halaman Instance list, klik nama instans tujuan.
-
Pada halaman Topics, temukan topik tujuan lalu pilih di kolom Actions.
-
Pada halaman Send and Consume Message with Quick Experience, masukkan Message Content.
Contoh ini mengirim pesan
1,adan2,b. -
Klik OK.
-
-
Login ke kluster ApsaraDB for ClickHouse, kueri tabel terdistribusi, dan periksa apakah data telah tersinkronisasi.
Untuk informasi lebih lanjut tentang cara login ke kluster ApsaraDB for ClickHouse, lihat Menghubungkan ke kluster ApsaraDB for ClickHouse menggunakan DMS.
Gunakan pernyataan berikut untuk mengkueri dan memverifikasi data:
Edisi Enterprise
SELECT * FROM kafka_table_local;Edisi yang kompatibel dengan Community
Kode berikut memberikan contoh cara mengkueri tabel terdistribusi.
-
Jika tabel tujuan adalah tabel lokal, Anda harus mengganti nama tabel terdistribusi dalam kueri dengan nama tabel lokal.
-
Jika Anda menggunakan kluster Edisi yang kompatibel dengan Community bertipe multi-node, kami sangat menyarankan agar Anda mengkueri tabel terdistribusi. Jika Anda langsung mengkueri tabel lokal, hanya data dari satu node yang akan dikembalikan, sehingga menghasilkan set hasil yang tidak lengkap.
SELECT * FROM kafka_table_distributed;Jika kueri mengembalikan hasil, maka sinkronisasi data dari Kafka ke ApsaraDB for ClickHouse berhasil.
Hasil kueri adalah sebagai berikut.
┌─id─┬─name─┐ │ 1 │ a │ │ 2 │ b │ └────┴──────┘Jika hasil kueri tidak sesuai harapan, lanjutkan ke Langkah 7 (Opsional): Periksa status konsumsi tabel eksternal Kafka untuk melakukan pemecahan masalah lebih lanjut.
-
Langkah 7 (Opsional): Periksa status konsumsi Kafka
Jika data yang tersinkronisasi tidak sesuai dengan data di Kafka, kueri tabel sistem untuk memeriksa status konsumsi tabel eksternal Kafka dan identifikasi pengecualian.
Engine v23.8 atau lebih baru
Jalankan pernyataan berikut untuk mengkueri tabel sistem system.kafka_consumers dan melihat status konsumsi tabel eksternal Kafka:
select * from system.kafka_consumers;
Tabel berikut menjelaskan bidang-bidang pada tabel system.kafka_consumers.
|
Field |
Description |
|
database |
Database tempat tabel eksternal Kafka berada. |
|
table |
Nama tabel eksternal Kafka. |
|
consumer_id |
ID konsumen Kafka. Satu tabel dapat memiliki beberapa konsumen. Jumlah konsumen ditentukan oleh parameter kafka_num_consumers saat Anda membuat tabel eksternal Kafka. |
|
assignments.topic |
Topik Kafka. |
|
assignments.partition_id |
ID partisi Kafka. Satu partisi hanya dapat ditugaskan ke satu konsumen. |
|
assignments.current_offset |
Offset saat ini. |
|
exceptions.time |
Timestamp dari 10 pengecualian terbaru. |
|
exceptions.text |
Teks dari 10 pengecualian terbaru. |
|
last_poll_time |
Timestamp polling terakhir. |
|
num_messages_read |
Jumlah pesan yang dibaca oleh konsumen. |
|
last_commit_time |
Timestamp commit terakhir. |
|
num_commits |
Total jumlah commit yang dilakukan oleh konsumen. |
|
last_rebalance_time |
Timestamp Penyeimbangan ulang Kafka terakhir. |
|
num_rebalance_revocations |
Jumlah kali partisi dicabut dari konsumen. |
|
num_rebalance_assignments |
Jumlah kali konsumen ditugaskan partisi dalam kluster Kafka. |
|
is_currently_used |
Menunjukkan apakah konsumen sedang digunakan. |
|
last_used |
Waktu terakhir konsumen digunakan, dalam Unix time (mikrodetik). |
|
rdkafka_stat |
Statistik internal library. Untuk informasi lebih lanjut, lihat librdkafka. Nilai default-nya adalah 3000, yang berarti statistik dihasilkan setiap 3 detik. Catatan
Ketika |
Engine sebelum v23.8
Jalankan pernyataan berikut untuk mengkueri tabel sistem system.kafka dan melihat status konsumsi tabel eksternal Kafka:
SELECT * FROM system.kafka;
Tabel berikut menjelaskan bidang-bidang pada tabel system.kafka.
|
Field |
Description |
|
database |
Nama database tempat tabel eksternal Kafka berada. |
|
table |
Nama tabel eksternal Kafka. |
|
topic |
Nama topik yang dikonsumsi oleh tabel eksternal Kafka. |
|
consumer_group |
Nama kelompok konsumen yang digunakan oleh tabel eksternal Kafka. |
|
last_read_message_count |
Jumlah pesan yang ditarik dari tabel eksternal Kafka. |
|
status |
Status konsumsi pesan Kafka oleh tabel eksternal. Nilai yang valid:
|
|
exception |
Detail pengecualian. Catatan
Jika nilai status adalah error, parameter ini mengembalikan detail pengecualian. |