All Products
Search
Document Center

ApsaraDB for ClickHouse:Sinkronisasi data dari Kafka

Last Updated:Aug 27, 2026

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

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.

Penting

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.

  • Jika Anda menggunakan ApsaraMQ for Kafka, ApsaraDB for ClickHouse dapat mengurai nama domain instans ApsaraMQ for Kafka secara default.

  • Jika Anda menggunakan cluster Kafka yang dikelola sendiri, ApsaraDB for ClickHouse mendukung koneksi ke cluster Kafka menggunakan alamat IP atau nama domain kustom dalam format tetap. Aturan nama domain kustom berikut didukung:

    1. Nama domain yang diakhiri dengan .com.

    2. Nama domain yang diakhiri dengan .local dan mengandung kafka, mysql, atau rabbitmq.

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
  1. Jika throughput satu konsumen tidak mencukupi, Anda harus menentukan lebih banyak konsumen.

  2. Jumlah total konsumen tidak boleh melebihi jumlah partisi dalam topik karena hanya satu konsumen yang dapat ditugaskan ke setiap partisi.

kafka_thread_per_consumer

Tidak

Menentukan apakah akan mengaktifkan thread khusus untuk setiap konsumen. Nilai default adalah 0. Nilai yang valid:

  1. 0: Semua konsumen berbagi satu thread untuk mengonsumsi data.

  2. 1: Thread khusus diaktifkan untuk setiap konsumen untuk mengonsumsi data.

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_skip_broken_messages=N, mesin akan melewati N pesan Kafka yang tidak dapat diurai. Satu pesan setara dengan satu baris data.

kafka_commit_every_batch

Tidak

Frekuensi commit Kafka. Nilai default adalah 0. Nilai yang valid:

  1. 0: Commit dilakukan hanya setelah blok data lengkap ditulis.

  2. 1: Commit dilakukan setelah setiap batch data ditulis.

kafka_auto_offset_reset

Tidak

Offset tempat pembacaan data Kafka dimulai. Nilai yang valid:

  1. earliest: membaca data Kafka dari offset paling awal. Ini adalah nilai default.

  2. latest: membaca data Kafka dari offset terbaru.

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.

Penting

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

  1. Buat tabel lokal.

    CREATE TABLE default.kafka_table_local ON CLUSTER default (
      id Int32,
      name String
    ) ENGINE = MergeTree()
    ORDER BY (id);
  2. (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

  1. Buat tabel lokal.

    CREATE TABLE default.kafka_table_local ON CLUSTER default (
      id Int32,
      name String
    ) ENGINE = ReplicatedMergeTree()
    ORDER BY (id);
  2. (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.

Penting

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.

  • Kluster Edisi Kompatibel Komunitas:

    • Untuk kluster multi-node, kami menyarankan Anda mengimpor data ke tabel terdistribusi.

    • Jika melakukan sinkronisasi ke tabel lokal, tentukan nama tabel lokal.

  • Kluster Edisi Perusahaan: Karena kluster Edisi Perusahaan tidak memiliki tabel terdistribusi, Anda harus menentukan tabel lokal.

  • Edisi Kompatibel Komunitas contoh: kafka_table_distributed

  • Edisi Perusahaan contoh: kafka_table_local

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

  1. Kirim pesan ke topik dalam instans ApsaraMQ for Kafka.

    1. Login ke Konsol ApsaraMQ for Kafka.

    2. Pada halaman Instance list, klik nama instans tujuan.

    3. Pada halaman Topics, temukan topik tujuan lalu pilih More > Send Message (Demo) di kolom Actions.

    4. Pada halaman Send and Consume Message with Quick Experience, masukkan Message Content.

      Contoh ini mengirim pesan 1,a dan 2,b.

    5. Klik OK.

  2. 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 statistics_interval_ms=0 dikonfigurasi untuk ApsaraDB for ClickHouse, pengumpulan statistik untuk tabel eksternal Kafka dinonaktifkan.

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:

  • no_view: Tidak ada tampilan yang dibuat untuk tabel eksternal Kafka.

  • attach_view: Tampilan telah dibuat untuk tabel eksternal Kafka.

  • normal: Status normal.

    Status normal menunjukkan bahwa tabel eksternal mengonsumsi data sesuai harapan.

  • skip_parse: Kesalahan parsing dilewati.

  • error: Terjadi pengecualian konsumsi.

exception

Detail pengecualian.

Catatan

Jika nilai status adalah error, parameter ini mengembalikan detail pengecualian.

FAQ