Topik ini menjelaskan cara membuat konektor sink Tablestore untuk mengekspor data dari topik ApsaraMQ for Kafka ke tabel Tablestore.
Prasyarat
-
Anda telah mengaktifkan Tablestore dan membuat instans. Untuk informasi selengkapnya, lihat Aktifkan Tablestore dan buat instans.
-
Anda telah membeli dan mengaktifkan instans ApsaraMQ for Kafka serta membuat topik. Untuk langkah-langkah detailnya, lihat Beli instans ApsaraMQ for Kafka dan Buat resource.
-
Peran terkait layanan yang dihasilkan oleh tugas konektor sink Tablestore memerlukan kebijakan
AliyunOTSFullAccess. Lampirkan kebijakan ini secara manual untuk memberikan izin kepada peran tersebut dalam mengelola Tablestore. Untuk informasi selengkapnya, lihat Berikan izin kepada Peran RAM.
Langkah 1: Buat tabel Tablestore
Buat tabel Tablestore untuk menyimpan data yang dialirkan dari ApsaraMQ for Kafka. Untuk informasi selengkapnya, lihat Prosedur.
Topik ini menggunakan contoh instans bernama ots-sink dan tabel data bernama ots_sink_table. Saat membuat tabel, definisikan tiga kunci primer: topic bertipe STRING (ditetapkan sebagai kunci partisi), partition bertipe INTEGER, dan offset bertipe INTEGER.
Langkah 2: Buat dan mulai konektor sink Tablestore
Masuk 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.
-
Pada halaman Create Task, atur Task Name dan Description, konfigurasikan parameter berikut, lalu klik Save.
-
Create 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 dari 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 merutekan 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 muatan menggunakan UTF-8.
-
Text: Mengkodekan data biner sebagai string UTF-8 ke dalam muatan. Ini adalah format default.
-
Binary: Mengkodekan data biner sebagai string Base64 ke dalam muatan.
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 mengumpulkan pesan dan mengirimkannya ke sink pada 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, pengayaan, dan perutean dinamis. Untuk informasi selengkapnya, lihat Gunakan Function Compute untuk membersihkan data pesan.
-
Pada langkah Sink, atur Service Type ke Tablestore dan konfigurasikan parameter berikut.
Parameter
Deskripsi
Contoh
Instance Name
Nama instans Tablestore tujuan.
ots-sink
Destination Table
Nama tabel data Tablestore tujuan.
ots_sink_table
Primary Key
Untuk menghasilkan nilai kunci primer Tablestore dan kolom atribut, Anda harus mendefinisikan aturan ekstraksi untuk konten setiap kolom atribut menggunakan sintaks JsonPath. Ketika Data Format pada langkah Source diatur ke JSON, data keluaran dari ApsaraMQ for Kafka memiliki format berikut:
{ "data": { "topic": "demo-topic", "partition": 0, "offset": 2, "timestamp": 1739756629123, "headers": { "headers": [], "isReadOnly": false }, "key":"ots-sink-k1", "value": "ots-sink-v1" }, "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" }Sebagai contoh, untuk kunci primer bernama topic, atur aturan ekstraksi nilainya menjadi
$.data.topic.Attribute Column
Sebagai contoh, untuk kolom atribut bernama key, atur aturan ekstraksi nilainya menjadi
$.data.key.Operation Mode
Metode penulisan data ke Tablestore.
-
put: Jika catatan dengan kunci primer yang sama sudah ada, data baru akan menimpa data yang ada.
-
update: Ketika dua catatan data memiliki kunci primer yang sama, kolom tambahan ditulis ke baris tersebut, dan kolom yang sudah ada tidak dihapus.
-
delete: Menghapus data kunci primer yang sesuai.
put
Network configuration
-
VPC: Gunakan VPC untuk mengirimkan pesan Kafka ke Tablestore.
-
Public Network: Mengirimkan pesan Kafka ke Tablestore melalui jaringan publik.
VPC
VPC
Pilih ID VPC. Parameter ini diperlukan hanya jika Network Configuration diatur ke VPC.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Pilih ID vSwitch. Parameter ini diperlukan hanya jika Network Configuration diatur ke VPC.
vsw-bp1gbjhj53hdjdkg****
Security group
Pilih security group. Parameter ini diperlukan hanya jika Network Configuration diatur ke VPC.
test_group
-
-
-
Task Properties
Konfigurasikan kebijakan pengulangan untuk pengiriman event yang gagal dan metode penanganan error. Untuk informasi selengkapnya, lihat Kebijakan pengulangan dan antrian surat mati.
-
-
Kembali ke halaman Tasks. Temukan tugas Anda dan klik Enable di kolom Actions.
-
Pada kotak dialog Note, baca informasinya dan klik OK.
Setelah diaktifkan, tugas akan mulai berjalan dalam waktu 30 hingga 60 detik. Anda dapat memantau statusnya di kolom Status pada halaman Tasks.
Langkah 3: Uji konektor sink Tablestore
-
Pada halaman Tasks, di kolom Event Source untuk tugas Anda, klik topik sumber.
- Pada halaman detail topik, klik Send Test Message.
-
Pada panel Start to Send and Consume Message, konfigurasikan pesan dan klik OK.
Pada tab Console, atur Message Key ke
ots-sink-k1, Message Content keots-sink-v1, dan Send to Specified Partition ke No. -
Kembali ke halaman Tasks, dan di kolom Event Target untuk tugas Anda, klik nama tabel tujuan.
-
Pada halaman Manage Table, klik tab Data Management untuk melihat data di tabel Tablestore.
Tabel data berisi kolom topic (kunci primer), partition (kunci primer), offset (kunci primer), key, dan value. Catatan dengan nilai seperti partition=
2, offset=10, key=ots-sink-k1, dan value=ots-sink-v1mengonfirmasi bahwa pesan Kafka telah ditulis ke Tablestore.