All Products
Search
Document Center

ApsaraMQ for Kafka:Buat konektor sink AnalyticDB

Last Updated:Jun 21, 2026

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.

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

  1. Login ke ApsaraMQ for Kafka console. Pada halaman Overview, pilih wilayah di bagian Resource Distribution.

  2. Di panel navigasi sebelah kiri, pilih Connector Ecosystem Integration > Tasks.

  3. Pada halaman Tasks, klik Create Task.

    • Konfigurasikan task

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

      2. Pada langkah Filtering, atur Pattern Content untuk memfilter event. Untuk informasi selengkapnya, lihat event pattern.

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

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

        Penting

        Setelah 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

  4. Kembali ke halaman Tasks. Temukan task yang Anda buat dan klik Enable di kolom Actions.

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

  1. Pada halaman Tasks, temukan task Anda dan klik nama topik sumber di kolom Event Source.

  2. Pada halaman detail topik, klik Send Test Message.
  3. Pada panel Start to Send and Consume Message, konfigurasikan isi pesan, lalu klik OK.

    Catatan

    Isi 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-k1 dan Message body ke {"userid":"xiaoming","source":"shanghai"}. Untuk Send to a specific partition, pilih No, lalu klik OK.

  4. Pada halaman Tasks, temukan task Anda dan klik nama instans tujuan di kolom Event Target.

  5. Pada halaman Basic Information instans, klik Log On to Database di pojok kanan atas.

  6. 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 userid bernilai xiaoming dan source bernilai shanghai. Hal ini mengonfirmasi bahwa data berhasil ditulis ke tabel tujuan.