Topik ini menjelaskan cara membuat konektor sink AnalyticDB untuk mengalirkan data dari topik sumber di instans ApsaraMQ for Kafka ke tabel dalam database AnalyticDB.
Prasyarat
Untuk informasi selengkapnya, lihat Prasyarat.
Langkah 1: Buat resource tujuan
Buat resource AnalyticDB for MySQL atau AnalyticDB for PostgreSQL.
-
AnalyticDB for MySQL: Di AnalyticDB for MySQL console, buat kluster dan akun database, hubungkan ke kluster, lalu buat database. Untuk informasi selengkapnya, lihat Buat kluster, Buat akun database, Hubungkan ke kluster, dan Buat database.
-
AnalyticDB for PostgreSQL: Di AnalyticDB for PostgreSQL console, buat instans dan akun database, lalu login ke database. Untuk informasi selengkapnya, lihat Buat instans, Buat dan kelola pengguna, dan Koneksi klien.
Contoh ini menggunakan database AnalyticDB for MySQL bernama adb_sink_database dan tabel bernama adb_sink_table.
Langkah 2: Buat dan aktifkan konektor sink AnalyticDB
Login ke ApsaraMQ for Kafka console. Pada halaman Overview, pilih wilayah di bagian Resource Distribution.
-
Di panel navigasi sebelah kiri, pilih .
-
Pada halaman Tasks, klik Create Task.
-
Konfigurasikan task
-
Pada langkah Source, pilih Message Queue for Apache Kafka sebagai Data Provider. Konfigurasikan parameter berikut, lalu klik Next.
Parameter
Deskripsi
Contoh
Region
Wilayah instans sumber Message Queue for Apache Kafka.
China (Beijing)
Kafka instance
Instans sumber Message Queue for Apache Kafka.
alikafka_post-cn-jte3****
Topic
Pilih topik untuk memproduksi pesan Message Queue for Apache Kafka.
demo-topic
Group ID
Kelompok konsumen instans sumber.
-
Quick Create: Opsi ini direkomendasikan. Group ID dengan format
GID_EVENTBRIDGE_xxxakan dibuat secara otomatis. -
Use Existing: Pilih group ID independen. Jangan gunakan group ID yang dibagi dengan layanan lain untuk menghindari gangguan pada konsumsi pesan yang sudah ada.
Quick Create
Consumer offset
Offset tempat konsumsi pesan dimulai.
-
Latest offset (latest)
-
Earliest offset (earliest)
Latest offset (latest)
Network configuration
Jenis jaringan untuk perutean pesan.
-
Basic Network
-
Self-managed Internet
Basic Network
VPC
Diperlukan hanya jika Network configuration diatur ke Self-managed Internet.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Diperlukan hanya jika Network configuration diatur ke Self-managed Internet.
vsw-bp1gbjhj53hdjdkg****
Security group
Diperlukan hanya jika Network configuration diatur ke Self-managed Internet.
alikafka_pre-cn-7mz2****
Data Format
Format encoding untuk konten pesan. Untuk sumber data yang mendukung transmisi biner, pilih salah satu format berikut. Kami merekomendasikan Json jika Anda tidak memiliki persyaratan encoding khusus.
-
Json: Mengkodekan data biner sebagai objek JSON ke dalam payload menggunakan UTF-8.
-
Text: Mengkodekan data biner sebagai string UTF-8 ke dalam payload. Ini adalah format default.
-
Binary: Mengkodekan data biner sebagai string Base64-encoded ke dalam payload.
Json
Messages
Parameter Advanced configuration. Jumlah maksimum pesan yang dikirim per batch. Permintaan dikirim hanya ketika jumlah pesan terakumulasi mencapai nilai yang ditentukan. Nilai valid: 1 hingga 10.000.
100
Interval (Unit: Seconds)
Parameter Advanced configuration. Interval pemanggilan fungsi. Sistem mengagregasi pesan dan mengirimkannya ke sink setiap interval. Nilai valid: 0 hingga 15. Nilai 0 berarti pesan dikirimkan segera tanpa penundaan.
3
-
-
Pada langkah Filtering, atur Pattern Content untuk memfilter event. Untuk informasi selengkapnya, lihat event pattern.
-
Pada langkah Transformation, konfigurasikan transformasi data untuk melakukan operasi seperti pemisahan, pemetaan, enrichmen, dan perutean dinamis. Untuk informasi selengkapnya, lihat Gunakan Function Compute untuk membersihkan data pesan.
-
Pada langkah Sink, atur Service Type ke AnalyticDB, konfigurasikan parameter berikut, lalu klik Save.
Parameter
Deskripsi
Contoh
Instance type
Pilih jenis database instans tujuan Anda.
-
AnalyticDB for MySQL
-
AnalyticDB for PostgreSQL
MySQL version
AnalyticDB instance ID
Pilih instans tujuan.
gp-bp10uo5n536wd****
Database name
Pilih database tujuan.
adb_sink_database
Table name
Pilih tabel data tujuan.
adb_sink_table
Data Mapping
Gunakan ekspresi JSONPath untuk menentukan aturan ekstraksi data. Ketika Data Format diatur ke Json pada langkah Source, data yang dialirkan dari ApsaraMQ for Kafka dibungkus dalam struktur CloudEvents, seperti yang ditunjukkan di bawah ini:
{ "data": { "topic": "demo-topic", "partition": 0, "offset": 2, "timestamp": 1739756629123, "headers": { "headers": [], "isReadOnly": false }, "key":"adb-sink-k1", "value": { "userid":"xiaoming", "source":"shanghai" } }, "id": "7702ca16-f944-4b08-***-***-0-2", "source": "acs:alikafka", "specversion": "1.0", "type": "alikafka:Topic:Message", "datacontenttype": "application/json; charset=utf-8", "time": "2025-02-17T01:43:49.123Z", "subject": "acs:alikafka:alikafka_serverless-cn-lf6418u6701:topic:demo-topic", "aliyunaccountid": "1******6789" }Petakan setiap kolom tabel tujuan ke bidang dalam pesan sumber menggunakan ekspresi JSONPath. Misalnya, untuk memetakan bidang
useriddari pesan ke kolom tabel, gunakan ekspresi$.data.value.userid.Database username
Masukkan username untuk akun database.
user
Database password
Masukkan password untuk akun database.
******
Network configuration
-
VPC: Hubungkan ke AnalyticDB melalui VPC.
-
Public Network: Hubungkan ke AnalyticDB melalui internet publik.
VPC
VPC
Pilih ID VPC. Parameter ini diperlukan hanya ketika Network configuration diatur ke VPC.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Pilih ID vSwitch. Parameter ini diperlukan hanya ketika Network configuration diatur ke VPC.
PentingSetelah Anda memilih vSwitch, Anda harus menambahkan Blok CIDR vSwitch ke daftar putih alamat IP instans AnalyticDB for MySQL. Untuk informasi selengkapnya, lihat Konfigurasikan daftar putih alamat IP.
vsw-bp1gbjhj53hdjdkg****
Security group
Pilih security group. Parameter ini diperlukan hanya ketika Network configuration diatur ke VPC.
test_group
-
-
-
-
Kembali ke halaman Tasks. Temukan task yang Anda buat dan klik Enable di kolom Actions.
-
Di kotak dialog Note, baca pesan tersebut lalu klik OK.
Task memerlukan waktu 30 hingga 60 detik untuk mulai berjalan. Anda dapat memantau perkembangannya di kolom Status pada halaman Tasks.
Langkah 3: Verifikasi konektor sink AnalyticDB
-
Pada halaman Tasks, temukan task Anda dan klik nama topik sumber di kolom Event Source.
- Pada halaman detail topik, klik Send Test Message.
-
Pada panel Start to Send and Consume Message, konfigurasikan isi pesan, lalu klik OK.
CatatanIsi pesan harus dalam format JSON. Bidang yang ditentukan dalam aturan pemetaan data Anda akan diekstraksi dan ditulis ke kolom yang sesuai di tabel tujuan.
Di kotak dialog Start to Send and Consume Message, pilih tab Console. Atur Message key ke
adb-sink-k1dan Message body ke{"userid":"xiaoming","source":"shanghai"}. Untuk Send to a specific partition, pilih No, lalu klik OK. -
Pada halaman Tasks, temukan task Anda dan klik nama instans tujuan di kolom Event Target.
-
Pada halaman Basic Information instans, klik Log On to Database di pojok kanan atas.
-
Di Konsol Data Management Service (DMS), jalankan Pernyataan SQL berikut untuk mengkueri semua data di tabel.
SELECT * FROM adb_sink_table;Kueri tersebut seharusnya mengembalikan catatan dengan
useridbernilaixiaomingdansourcebernilaishanghai. Hal ini mengonfirmasi bahwa data berhasil ditulis ke tabel tujuan.