AnalyticDB for MySQL menyediakan fitur sinkronisasi data melalui AnalyticDB Pipeline Service (APS). Anda dapat membuat tugas sinkronisasi Kafka untuk mengingesti data dari Kafka secara real time mulai dari offset tertentu, sehingga menghasilkan output data hampir real time, arsip data historis lengkap, dan analitik elastis. Topik ini menjelaskan cara menambahkan sumber data Kafka, membuat dan menjalankan tugas sinkronisasi Kafka, serta melakukan analisis data dan manajemen sumber data setelah data disinkronkan.
Prasyarat
-
Kluster AnalyticDB for MySQL Edisi Perusahaan, Edisi Dasar, atau Edisi Data Lakehouse telah dibuat.
Akun database telah dibuat untuk kluster AnalyticDB for MySQL.
Jika Anda menggunakan akun Alibaba Cloud, Anda hanya perlu membuat akun istimewa.
Jika Anda menggunakan pengguna Resource Access Management (RAM), Anda harus membuat akun istimewa dan akun standar serta mengaitkan akun standar tersebut dengan pengguna RAM.
-
Anda telah membuat instans ApsaraMQ for Kafka (Kafka) di wilayah yang sama dengan kluster AnalyticDB for MySQL.
-
Anda telah membuat topik Kafka dan mengirimkan pesan ke dalamnya. Untuk informasi selengkapnya, lihat Panduan cepat Message Queue for Apache Kafka.
Catatan
-
Hanya data Kafka dalam format JSON yang didukung.
-
Data dalam topik Kafka akan dihapus secara otomatis ketika periode retensinya berakhir. Jika pekerjaan sinkronisasi data gagal, data mungkin terhapus dari topik sebelum pekerjaan dimulai ulang, sehingga menyebabkan kehilangan data. Untuk mencegah hal ini, tingkatkan periode retensi data topik dan segera hubungi dukungan teknis jika pekerjaan gagal.
-
API Kafka memotong data sampel dari Kafka yang melebihi 8 KB. Hal ini menyebabkan penguraian JSON gagal, sehingga sistem tidak dapat menghasilkan pemetaan bidang secara otomatis.
-
Perubahan pada skema tabel sumber Kafka tidak memicu pembaruan DDL otomatis dan tidak disinkronkan ke AnalyticDB for MySQL.
-
Setelah ingesti data, operasi commit diperlukan agar data menjadi terlihat. Untuk memastikan stabilitas pekerjaan dan kinerja baca/tulis optimal, fitur sinkronisasi data AnalyticDB for MySQL menggunakan interval commit default 5 menit. Oleh karena itu, setelah Anda membuat dan menjalankan pekerjaan sinkronisasi data, Anda harus menunggu minimal 5 menit untuk melihat batch data pertama.
Penagihan
Biaya berikut berlaku saat Anda menggunakan fitur migrasi data AnalyticDB for MySQL untuk memigrasikan data ke Object Storage Service (OSS):
-
Biaya sumber daya ACU elastis AnalyticDB for MySQL. Untuk detailnya, lihat Penagihan untuk Edisi Data Lakehouse dan Penagihan untuk Edisi Perusahaan dan Edisi Dasar.
-
Biaya penyimpanan OSS, biaya permintaan GET, serta biaya permintaan PUT dan lainnya. Untuk detailnya, lihat Ikhtisar penagihan.
Alur Kerja
-
Langkah 1: Buat sumber data.
-
Langkah 2: Buat tautan sinkronisasi.
-
Langkah 3: Jalankan tugas sinkronisasi data.
-
Langkah 4: Analisis data.
-
Langkah 5 (Opsional): Kelola sumber data.
Buat sumber data
Jika Anda telah membuat sumber data Kafka, Anda dapat melewati langkah ini dan langsung membuat tautan data. Untuk informasi selengkapnya, lihat Buat tautan data.
Masuk ke Konsol AnalyticDB for MySQL. Di pojok kiri atas konsol, pilih wilayah. Di panel navigasi sebelah kiri, klik Clusters. Temukan kluster yang ingin Anda kelola lalu klik ID kluster tersebut.
-
Di panel navigasi sebelah kiri, pilih Data Ingestion>Data Sources.
-
Di pojok kiri atas, klik Create Data Source.
-
Pada halaman Create Data Source, konfigurasikan parameter berikut.
Parameter
Deskripsi
Data Source Type
Pilih Kafka.
Data Source Name
Sistem menghasilkan nama berdasarkan jenis sumber data dan waktu saat ini. Anda dapat mengubah nama sesuai kebutuhan.
Data Source Description
Deskripsi sumber data, seperti kasus penggunaannya atau batasan bisnisnya.
Deployment Mode
Saat ini, hanya Alibaba Cloud Instance yang didukung.
Kafka Instance
ID instans Kafka.
Masuk ke Konsol ApsaraMQ for Kafka. Di halaman Instances, lihat ID instans.
Kafka Topic
Nama topik Kafka Anda.
Masuk ke Konsol ApsaraMQ for Kafka. Di halaman Topics instans target, lihat nama topik.
Message Data Format
Format data untuk pesan. Saat ini, hanya JSON yang didukung.
-
Setelah mengonfigurasi parameter, klik Create.
Buat pekerjaan sinkronisasi
-
Di panel navigasi sebelah kiri, klik Simple Log Service/Kafka Data Synchronization.
-
Di pojok kiri atas, klik Create Synchronization Job.
-
Pada halaman Create Synchronization Job, konfigurasikan parameter pada bagian Source and Destination Settings, Destination Database and Table Settings, dan Synchronization Settings.
-
Tabel berikut menjelaskan parameter untuk Source and Destination Settings:
Parameter
Deskripsi
Job Name
Masukkan nama untuk pekerjaan sinkronisasi. Secara default, sistem menghasilkan nama berdasarkan sumber data dan waktu saat ini, yang dapat Anda sesuaikan.
Data Source
Pilih sumber data Kafka yang sudah ada, atau buat yang baru.
Destination Type
Pilih salah satu opsi berikut:
-
Data Lake - User OSS.
-
Data Lake - AnalyticDB Lake Storage (disarankan).
PentingJika Anda memilih Data Lake - AnalyticDB Lake Storage, Anda harus terlebih dahulu mengaktifkan fitur penyimpanan lake.
ADB Lake Storage
Nama penyimpanan lake untuk data lake AnalyticDB for MySQL.
Pilih penyimpanan lake tujuan dari daftar drop-down. Jika tidak tersedia penyimpanan lake, klik Automatically Created untuk membuatnya.
PentingParameter ini wajib diisi hanya jika Destination Type diatur ke Data Lake - AnalyticDB Lake Storage.
OSS Path
Jalur penyimpanan data lake AnalyticDB for MySQL di OSS.
Penting-
Parameter ini wajib diisi hanya jika Destination Type diatur ke Data Lake - User OSS.
-
Daftar drop-down menampilkan semua bucket di wilayah yang sama dengan kluster AnalyticDB for MySQL. Anda dapat memilih salah satunya. Rencanakan jalur penyimpanan dengan cermat. Anda tidak dapat mengubah jalur ini setelah pekerjaan dibuat.
-
Kami menyarankan Anda memilih direktori kosong. Jalur OSS tidak boleh merupakan awalan dari jalur OSS pekerjaan sinkronisasi lain, atau sebaliknya. Hal ini mencegah data saling menimpa. Misalnya, jika dua pekerjaan sinkronisasi memiliki jalur OSS
oss://testBucketName/test/sls1/danoss://testBucketName/test/, kedua jalur tersebut memiliki hubungan awalan, sehingga data mungkin saling menimpa selama sinkronisasi.
Storage Format
Format penyimpanan data. Opsi yang didukung meliputi:
-
PAIMON.
PentingFormat ini hanya didukung jika Destination Type diatur ke Data Lake - User OSS.
-
ICEBERG.
-
-
Tabel berikut menjelaskan parameter untuk Destination Database and Table Settings:
Parameter
Deskripsi
Database Name
Nama database tujuan di AnalyticDB for MySQL. Jika database dengan nama tersebut belum ada, database baru akan dibuat. Jika tidak, data akan disinkronkan ke database yang sudah ada. Untuk informasi tentang konvensi penamaan, lihat Batasan.
PentingDi bagian Source and Destination Settings, jika Storage Format diatur ke PAIMON, database yang sudah ada harus memenuhi kondisi berikut. Jika tidak, pekerjaan sinkronisasi akan gagal:
-
Database harus merupakan database eksternal. Pernyataan
CREATE DATABASEharus berupaCREATE EXTERNAL DATABASE <database_name>. -
Klausa
DBPROPERTIESdalam pernyataanCREATE DATABASEharus mencakup properticatalog, dan nilaicatalogharuspaimon. -
Klausa
DBPROPERTIESharus mencakup propertiadb.paimon.warehouse. Contoh:adb.paimon.warehouse=oss://testBucketName/aps/data. -
Klausa
DBPROPERTIESharus mencakup propertiLOCATION, dan Anda harus menambahkan.dbke nama database dalam jalur tersebut. Jika tidak, kueri XIHE akan gagal. Contoh:LOCATION='oss://testBucketName/aps/data/kafka_paimon_external_db.db/'.Untuk jalur OSS yang ditentukan oleh
LOCATION, bucket dan direktori harus sudah ada. Jika tidak, pembuatan database akan gagal.
Table Name
Nama tabel tujuan di AnalyticDB for MySQL. Jika tabel dengan nama tersebut belum ada, tabel baru akan dibuat. Pekerjaan akan gagal jika tabel dengan nama tersebut sudah ada. Untuk informasi tentang konvensi penamaan, lihat Batasan.
Sample Data
Sistem secara otomatis mengambil data terbaru dari topik Kafka untuk digunakan sebagai data sampel.
CatatanData dalam topik Kafka harus dalam format JSON. Jika data dalam format lain, kesalahan akan terjadi selama sinkronisasi data.
Parsed JSON Layers
Tentukan jumlah level bersarang JSON yang akan diurai. Nilai yang valid:
-
0: Tidak mengurai.
-
1 (default): Mengurai satu level.
-
2: Mengurai dua level.
-
3: Mengurai tiga level.
-
4: Mengurai empat level.
Untuk informasi lebih lanjut tentang strategi penguraian JSON, lihat Contoh level penguraian JSON dan inferensi skema.
Schema Field Mapping
Bagian ini menampilkan skema yang diinferensi dari data sampel. Anda kemudian dapat mengubah nama bidang dan tipe data tujuan, atau menambah dan menghapus bidang.
Partition Key Settings
Tentukan kunci partisi untuk tabel tujuan. Kami menyarankan Anda mempartisi berdasarkan waktu log atau logika bisnis untuk meningkatkan kinerja ingesti dan kueri data. Jika Anda tidak mengatur parameter ini, tabel tujuan tidak akan dipartisi secara default.
Anda dapat memformat kunci partisi tujuan berdasarkan waktu atau berdasarkan bidang partisi tertentu.
-
Untuk mempartisi data berdasarkan tanggal dan waktu, pilih bidang tanggal/waktu sebagai bidang kunci partisi. Untuk metode pemformatan, pilih pemformatan waktu, lalu tentukan format bidang sumber dan format partisi tujuan. AnalyticDB for MySQL menggunakan format bidang sumber untuk mengidentifikasi nilai bidang tersebut, lalu mengonversinya ke format partisi tujuan. Misalnya, jika bidang sumber adalah
gmt_createddengan nilai1711358834, format bidang sumber adalah timestamp dengan presisi detik, dan format partisi tujuan adalahyyyyMMdd, data akan dipartisi berdasarkan20240325. -
Untuk mempartisi data berdasarkan nilai bidang, pilih Specify partition field sebagai metode pemformatan.
-
-
Tabel berikut menjelaskan parameter untuk Synchronization Settings:
Parameter
Deskripsi
Starting Consumer Offset for Incremental Synchronization
Menentukan titik awal konsumsi data dari Kafka oleh pekerjaan. Opsi yang valid:
-
Earliest offset (begin_cursor): Mengonsumsi data mulai dari titik waktu paling awal yang tersedia di Kafka.
-
Latest offset (end_cursor): Mengonsumsi data mulai dari titik waktu paling akhir yang tersedia di Kafka.
-
Custom offset: Pilih titik waktu. Sistem mulai mengonsumsi data dari catatan pertama yang timestamp-nya lebih besar dari atau sama dengan titik waktu yang dipilih.
Job Resource Group
Pilih Kelompok Sumber Daya Pekerjaan tempat pekerjaan dijalankan.
ACUs for Incremental Synchronization
Tentukan jumlah AnalyticDB Compute Units (ACUs) dalam Kelompok Sumber Daya Pekerjaan untuk pekerjaan tersebut. Nilai minimum adalah 2. Nilai maksimum adalah jumlah sumber daya komputasi yang tersedia di Kelompok Sumber Daya Pekerjaan. Kami menyarankan mengalokasikan lebih banyak ACUs untuk meningkatkan kinerja ingesti dan stabilitas pekerjaan.
CatatanSaat Anda membuat pekerjaan sinkronisasi data, pekerjaan tersebut menggunakan sumber daya elastis di Kelompok Sumber Daya Pekerjaan. Pekerjaan sinkronisasi data menempati sumber daya dalam jangka waktu lama, sehingga sistem mengurangi sumber daya yang digunakan oleh pekerjaan tersebut dari Kelompok Sumber Daya Pekerjaan. Misalnya, jika Kelompok Sumber Daya Pekerjaan memiliki maksimum 48 ACUs dan pekerjaan sinkronisasi yang sudah ada menggunakan 8 ACUs, jumlah maksimum ACUs yang tersedia untuk pekerjaan lain dalam Kelompok Sumber Daya Pekerjaan yang sama adalah 40.
Advanced Settings
Menyediakan opsi lanjutan untuk menyesuaikan pekerjaan sinkronisasi. Penerapan pengaturan ini memerlukan bantuan dari dukungan teknis.
-
-
-
Setelah mengonfigurasi parameter, klik Submit.
Jalankan tugas sinkronisasi data
-
Di halaman Simple Log Service/Kafka Data Synchronization, pilih tugas sinkronisasi data tersebut lalu klik Start di kolom Actions.
-
Klik Search di pojok kiri atas. Tugas berhasil dijalankan jika statusnya berubah menjadi Running.
Analisis data
Setelah tugas sinkronisasi berhasil selesai, Anda dapat menggunakan pengembangan Spark Jar untuk menganalisis data di AnalyticDB for MySQL. Untuk informasi lebih lanjut tentang pengembangan Spark, lihat Editor pengembangan Spark dan Pengembangan aplikasi offline Spark.
-
Di panel navigasi sebelah kiri, klik .
-
Di templat default, masukkan pernyataan contoh lalu klik Run Now.
-- Berikut ini hanya contoh SparkSQL. Ubah kontennya dan jalankan program spark Anda. conf spark.driver.resourceSpec=medium; conf spark.executor.instances=2; conf spark.executor.resourceSpec=medium; conf spark.app.name=Spark SQL Test; conf spark.adb.connectors=oss; -- Berikut adalah pernyataan sql Anda show tables from lakehouse20220413156_adbTest; -
Opsi: Di tab Applications, klik Logs di kolom Actions untuk melihat log eksekusi pekerjaan Spark SQL.
Kelola sumber data
Di panel navigasi sebelah kiri, klik Data Ingestion>Data Sources. Anda dapat melakukan tindakan berikut di kolom Actions.
|
Tindakan |
Deskripsi |
|
Create Job |
Membuat pekerjaan sinkronisasi data atau migrasi data untuk sumber data tersebut. |
|
View |
Menampilkan konfigurasi detail sumber data. |
|
Edit |
Memungkinkan Anda mengedit properti sumber data, seperti nama dan deskripsinya. |
|
Delete |
Menghapus sumber data. Catatan
Anda tidak dapat menghapus sumber data yang memiliki pekerjaan sinkronisasi atau migrasi data terkait. Anda harus terlebih dahulu menghapus pekerjaan tersebut di halaman Simple Log Service/Kafka Data Synchronization. Untuk melakukannya, temukan pekerjaan target, lalu di kolom Actions, klik Delete. |
Contoh penguraian JSON dan inferensi skema
Pengaturan Parsed JSON Layers mengontrol berapa banyak level bersarang objek JSON yang diratakan menjadi bidang terpisah. Sebagai contoh, pertimbangkan data JSON berikut yang dikirim ke Kafka:
{
"name" : "zhangle",
"age" : 18,
"device" : {
"os" : {
"test": "lag",
"member":{
"fa": "zhangsan",
"mo": "limei"
}
},
"brand" : "none",
"version" : "11.4.2"
}
}
Contoh-contoh berikut menunjukkan hasil penguraian untuk level 0 hingga 4.
Penguraian Level 0
Tidak dilakukan penguraian. Seluruh objek JSON dikeluarkan sebagai satu bidang.
|
Bidang JSON |
Nilai |
Nama tujuan |
|
__value__ |
{ "name" : "zhangle","age" : 18, "device" : { "os" : { "test":"lag","member":{ "fa":"zhangsan","mo":"limei" }},"brand": "none","version" : "11.4.2" }} |
__value__ |
Penguraian Level 1
Level pertama bidang JSON diurai.
|
Bidang JSON |
Nilai |
Nama tujuan |
|
name |
zhangle |
name |
|
age |
18 |
age |
|
device |
{ "os" : { "test":"lag","member":{ "fa":"zhangsan","mo":"limei" }},"brand": "none","version" : "11.4.2" } |
device |
Penguraian Level 2
Bidang non-bersarang, seperti name dan age, dikeluarkan secara langsung. Bidang bersarang diratakan menjadi sub-bidangnya. Misalnya, bidang bersarang device diratakan menjadi device.os, device.brand, dan device.version.
Karena nama bidang tujuan tidak mendukung titik (.), sistem secara otomatis menggantinya dengan garis bawah (_).
|
Bidang JSON |
Nilai |
Nama tujuan |
|
name |
zhangle |
name |
|
age |
18 |
age |
|
device.os |
{ "test":"lag","member":{ "fa":"zhangsan","mo":"limei" }} |
device_os |
|
device.brand |
none |
device_brand |
|
device.version |
11.4.2 |
device_version |
Penguraian Level 3
|
Bidang JSON |
Nilai |
Nama tujuan |
|
name |
zhangle |
name |
|
age |
18 |
age |
|
device.os.test |
lag |
device_os_test |
|
device.os.member |
{ "fa":"zhangsan","mo":"limei" } |
device_os_member |
|
device.brand |
none |
device_brand |
|
device.version |
11.4.2 |
device_version |
Penguraian Level 4
|
JSON Field |
Nilai |
Nama tujuan |
|
name |
zhangle |
name |
|
age |
18 |
age |
|
device.os.test |
lag |
device_os_test |
|
device.os.member.fa |
zhangsan |
device_os_member_fa |
|
device.os.member.mo |
limei |
device_os_member_mo |
|
device.brand |
none |
device_brand |
|
device.version |
11.4.2 |
device_version |