All Products
Search
Document Center

ApsaraMQ for Kafka:Buat konektor sink Tablestore

Last Updated:Jun 21, 2026

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

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

  4. Pada halaman Create Task, atur Task Name dan Description, konfigurasikan parameter berikut, lalu klik Save.

    • Create 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 dari 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 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

      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, pengayaan, dan perutean dinamis. Untuk informasi selengkapnya, lihat Gunakan Function Compute untuk membersihkan data pesan.

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

  5. Kembali ke halaman Tasks. Temukan tugas Anda dan klik Enable di kolom Actions.

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

  1. Pada halaman Tasks, di kolom Event Source untuk tugas Anda, klik topik sumber.

  2. Pada halaman detail topik, klik Send Test Message.
  3. 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 ke ots-sink-v1, dan Send to Specified Partition ke No.

  4. Kembali ke halaman Tasks, dan di kolom Event Target untuk tugas Anda, klik nama tabel tujuan.

  5. 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-v1 mengonfirmasi bahwa pesan Kafka telah ditulis ke Tablestore.