All Products
Search
Document Center

AnalyticDB:Sinkronkan data Kafka dengan APS (disarankan)

Last Updated:Aug 25, 2026

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

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):

Alur Kerja

Buat sumber data

Catatan

Jika Anda telah membuat sumber data Kafka, Anda dapat melewati langkah ini dan langsung membuat tautan data. Untuk informasi selengkapnya, lihat Buat tautan data.

  1. 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.

  2. Di panel navigasi sebelah kiri, pilih Data Ingestion>Data Sources.

  3. Di pojok kiri atas, klik Create Data Source.

  4. 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.

  5. Setelah mengonfigurasi parameter, klik Create.

Buat pekerjaan sinkronisasi

  1. Di panel navigasi sebelah kiri, klik Simple Log Service/Kafka Data Synchronization.

  2. Di pojok kiri atas, klik Create Synchronization Job.

  3. 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).

        Penting

        Jika 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.

      Penting

      Parameter 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/ dan oss://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.

        Penting

        Format 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.

      Penting

      Di 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 DATABASE harus berupa CREATE EXTERNAL DATABASE <database_name>.

      • Klausa DBPROPERTIES dalam pernyataan CREATE DATABASE harus mencakup properti catalog, dan nilai catalog harus paimon.

      • Klausa DBPROPERTIES harus mencakup properti adb.paimon.warehouse. Contoh: adb.paimon.warehouse=oss://testBucketName/aps/data.

      • Klausa DBPROPERTIES harus mencakup properti LOCATION, dan Anda harus menambahkan .db ke 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.

      Catatan

      Data 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_created dengan nilai 1711358834, format bidang sumber adalah timestamp dengan presisi detik, dan format partisi tujuan adalah yyyyMMdd, data akan dipartisi berdasarkan 20240325.

      • 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.

      Catatan

      Saat 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.

  4. Setelah mengonfigurasi parameter, klik Submit.

Jalankan tugas sinkronisasi data

  1. Di halaman Simple Log Service/Kafka Data Synchronization, pilih tugas sinkronisasi data tersebut lalu klik Start di kolom Actions.

  2. 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.

  1. Di panel navigasi sebelah kiri, klik Job Development > Spark JAR Development.

  2. 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;
  3. 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.

Penting

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