Topik ini menjelaskan cara membuat konektor sink MaxCompute untuk mengekspor data dari sebuah topik di instans ApsaraMQ for Kafka ke tabel MaxCompute.
Prasyarat
Untuk langkah-langkah detail, lihat Prasyarat untuk konektor sink.
Catatan
Untuk menggunakan fitur partisi MaxCompute, Anda harus membuat kolom partisi tambahan bernama time dengan tipe data STRING saat membuat tabel.
Langkah 1: Buat resource tujuan
Buat tabel menggunakan client MaxCompute. Untuk informasi selengkapnya, lihat Buat tabel.
Tutorial ini menggunakan tabel bernama kafka_to_maxcompute sebagai contoh. Tabel tersebut berisi tiga kolom dan menggunakan fitur partisi. Pernyataan SQL berikut membuat tabel tersebut:
CREATE TABLE IF NOT EXISTS kafka_to_maxcompute(topic STRING,valueName STRING,valueAge BIGINT) PARTITIONED by (time STRING);
Jika Anda tidak menggunakan fitur partisi, gunakan pernyataan berikut:
CREATE TABLE IF NOT EXISTS kafka_to_maxcompute(topic STRING,valueName STRING,valueAge BIGINT);
Setelah pernyataan berhasil dijalankan, hasil berikut dikembalikan:
Pada halaman Tables, informasi tabel kafka_to_maxcompute yang telah dibuat dapat dilihat. Jenis partisinya adalah partitioned table dan jenis tabelnya adalah internal table. Skema tabel mencakup tiga bidang: topic (string), valueName (string), dan valueAge (bigint). Tidak ada di antaranya yang merupakan primary key. Bidang partisinya adalah time (string).
Langkah 2: Buat dan mulai konektor
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, konfigurasikan parameter Task Name dan Description. Kemudian, ikuti petunjuk di layar untuk mengonfigurasi parameter lainnya.
-
Pembuatan Tugas
-
Pada wizard konfigurasi Source, atur Data Provider menjadi ApsaraMQ for Kafka, konfigurasikan parameter berikut, lalu klik Next.
Parameter
Deskripsi
Contoh
Region
Wilayah tempat instans ApsaraMQ for Kafka berada.
China (Hangzhou)
ApsaraMQ for Kafka Instance
ID instans ApsaraMQ for Kafka tempat data yang ingin Anda rutekan diproduksi.
alikafka_post-cn-9hdsbdhd****
Topic
Topik pada instans ApsaraMQ for Kafka tempat data yang ingin Anda rutekan diproduksi.
guide-sink-topic
Group ID
ID kelompok pada instans ApsaraMQ for Kafka tempat data yang ingin Anda rutekan diproduksi.
Quickly Create: Sistem secara otomatis membuat kelompok dengan ID dalam format GID_EVENTBRIDGE_xxx.
Use Existing Group: Pilih ID kelompok yang sudah ada dan tidak sedang digunakan. Jika Anda memilih kelompok yang sedang digunakan, publikasi dan langganan pesan yang ada akan terpengaruh.
Use Existing Group
Consumer Offset
Latest Offset: Pesan dikonsumsi dari offset terbaru.
Earliest Offset: Pesan dikonsumsi dari offset paling awal.
Latest Offset
Network Configuration
Jika diperlukan transmisi data lintas batas negara, pilih Self-managed Internet. Dalam kasus lain, pilih Basic Network.
Basic Network
Data Format
Fitur format data digunakan untuk mengencode data biner yang dikirimkan dari sumber ke format data tertentu. Beberapa format data didukung. Jika Anda tidak memiliki persyaratan khusus terkait encoding, tentukan Json sebagai nilainya.
Json: Mengencode data biner ke format JSON berdasarkan UTF-8 dan menempatkannya di payload.
Text: Format default. Mengencode data biner ke string berdasarkan UTF-8 dan menempatkannya di payload.
Binary: Data biner diencode ke string berdasarkan encoding Base64 lalu dimasukkan ke payload.
Text
Messages
Parameter Advanced Configuration. Jumlah maksimum pesan yang dikirim dalam satu batch. Permintaan hanya dikirim ketika jumlah pesan yang terakumulasi mencapai nilai ini. Nilainya harus berupa bilangan bulat dari 1 hingga 10.000.
2000
Interval (Unit: Seconds)
Parameter Advanced Configuration. Interval pemanggilan fungsi. Sistem mengagregasi pesan dan mengirimkannya ke Function Compute pada interval ini. Nilainya harus berupa bilangan bulat dari 0 hingga 15. Satuannya detik. Nilai 0 berarti pesan dikirimkan segera.
3
-
Pada langkah Filtering, definisikan pola event untuk menyaring permintaan. Untuk informasi selengkapnya, lihat Event patterns.
-
Pada langkah Transformation, konfigurasikan pembersihan data untuk menerapkan kemampuan pemrosesan data kompleks seperti pemisahan, pemetaan, pengayaan, dan perutean dinamis. Untuk informasi selengkapnya, lihat Gunakan Function Compute untuk membersihkan data pesan.
-
Pada langkah Sink, atur Service Type menjadi acs.maxcompute dan konfigurasikan parameter berikut.
Parameter
Deskripsi
Contoh
AccessKey ID
ID AccessKey untuk akun Alibaba Cloud Anda, yang digunakan untuk mengakses layanan MaxCompute.
yourAccessKeyID
AccessKey Secret
Rahasia AccessKey untuk akun Alibaba Cloud Anda.
yourAccessKeySecret
MaxCompute Project Name
Pilih proyek MaxCompute yang sudah ada.
test_compute
MaxCompute Table Name
Pilih tabel MaxCompute yang sudah ada.
kafka_to_maxcompute
MaxCompute Table Input Parameter
Setelah Anda memilih tabel, nama kolom dan tipe datanya akan ditampilkan. Anda hanya perlu mengonfigurasi value extraction rule untuk setiap kolom. Kode berikut menunjukkan contoh pesan. Dalam contoh ini, nilai untuk kolom
topicdiekstraksi dari bidangtopic. Oleh karena itu, value extraction rule didefinisikan sebagai$.topic.{ 'data': { 'topic': 't_test', 'partition': 2, 'offset': 1, 'timestamp': 1717048990499, 'headers': { 'headers': [], 'isReadOnly': False }, 'key': 'MaxCompute-K1', 'value': 'MaxCompute-V1' }, 'id': '9b05fc19-9838-4990-bb49-ddb942307d3f-2-1', 'source': 'acs:alikafka', 'specversion': '1.0', 'type': 'alikafka:Topic:Message', 'datacontenttype': 'application/json; charset=utf-8', 'time': '2024-05-30T06:03:10.499Z', 'aliyunaccountid': '1413397765616316' }topic:
$.data.topicvaluename:
$.data.valuevalueage:
$.data.offsetPartition Dimension
-
Close: Menonaktifkan fitur partisi.
-
Enable: Mengaktifkan fitur partisi.
Jika Anda mengaktifkan partisi, Anda harus mengonfigurasi parameter seperti nilai partisi:
-
Nilai partisi mendukung variabel waktu {yyyy}, {MM}, {dd}, {HH}, dan {mm}, yang masing-masing merepresentasikan tahun, bulan, hari, jam, dan menit. Variabel waktu bersifat case-sensitive.
-
Nilai partisi juga dapat berupa konstanta.
-
Enable
{yyyy}-{MM}-{dd}.{HH}:{mm}.suffix
Network configuration
-
VPC: Mengirimkan pesan Kafka ke MaxCompute melalui VPC.
-
Public Network: Mengirimkan pesan Kafka ke MaxCompute melalui internet.
internet
VPC
Pilih ID VPC. Parameter ini wajib hanya jika Anda mengatur Network Configuration menjadi VPC.
vpc-bp17fapfdj0dwzjkd****
vSwitch
Pilih ID vSwitch. Parameter ini wajib hanya jika Anda mengatur Network Configuration menjadi VPC.
vsw-bp1gbjhj53hdjdkg****
Security Group
Pilih grup keamanan. Parameter ini wajib hanya jika Anda mengatur Network Configuration menjadi VPC.
test_group
-
-
-
Properti Tugas
Konfigurasikan kebijakan retry dan dead-letter queue untuk menangani error pengiriman. Untuk informasi selengkapnya, lihat Kebijakan retry dan dead-letter queue.
-
-
Setelah menyelesaikan konfigurasi di atas, klik Save. Pada halaman Tasks, temukan tugas konektor sink MaxCompute yang telah Anda buat. Kolom Status menampilkan Starting. Ketika status berubah menjadi Running, proses pembuatan selesai.
Langkah 3: Uji konektor
-
Pada halaman Tasks, temukan konektor sink MaxCompute dan klik topik sumber di kolom Event Source.
- Pada halaman detail topik, klik Send Test Message.
-
Pada panel Start to Send and Consume Message, konfigurasikan konten pesan sebagai berikut, lalu klik OK.
Pada tab Console, atur Message key menjadi
MaxCompute-K1, atur Message body menjadiMaxCompute-V1, dan atur Send to specified partition menjadi No. -
Buka konsol MaxCompute dan jalankan pernyataan SQL berikut untuk melihat informasi partisi.
show PARTITIONS kafka_to_maxcompute;Hasil berikut dikembalikan:
OK OK OK OK OK OK OK OK OK OK OK OK time=2024-05-31.16:37.suffix OK 2024-05-31 16:42:49 INFO ================================================================== 2024-05-31 16:42:49 INFO Exit code of the Shell command 0 2024-05-31 16:42:49 INFO --- Invocation of Shell command completed --- 2024-05-31 16:42:49 INFO Shell run successfully! 2024-05-31 16:42:49 INFO Current task status: FINISH 2024-05-31 16:42:49 INFO Cost time is: 1.411s -
Berdasarkan informasi partisi, jalankan pernyataan berikut untuk melihat data dalam partisi tersebut.
SELECT * FROM kafka_to_maxcompute WHERE time="2024-05-31.16:37.suffix";Kueri mengembalikan satu catatan data dengan nilai kolom sebagai berikut:
topicadalahxxx(disembunyikan),valueNameadalahMaxCompute-V1,valueAgeadalah4, dantimeadalah2024-05-31.16:37.suffix. Hal ini menunjukkan bahwa data berhasil ditulis dari ApsaraMQ for Kafka ke tabel partisi di MaxCompute.