All Products
Search
Document Center

ApsaraMQ for Kafka:Buat konektor sink MaxCompute

Last Updated:Jun 22, 2026

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

  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, konfigurasikan parameter Task Name dan Description. Kemudian, ikuti petunjuk di layar untuk mengonfigurasi parameter lainnya.

    • Pembuatan Tugas

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

      2. Pada langkah Filtering, definisikan pola event untuk menyaring permintaan. Untuk informasi selengkapnya, lihat Event patterns.

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

      4. 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 topic diekstraksi dari bidang topic. 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.topic

        valuename: $.data.value

        valueage: $.data.offset

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

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

  1. Pada halaman Tasks, temukan konektor sink MaxCompute dan klik 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 konten pesan sebagai berikut, lalu klik OK.

    Pada tab Console, atur Message key menjadi MaxCompute-K1, atur Message body menjadi MaxCompute-V1, dan atur Send to specified partition menjadi No.

  4. 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
  5. 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: topic adalah xxx (disembunyikan), valueName adalah MaxCompute-V1, valueAge adalah 4, dan time adalah 2024-05-31.16:37.suffix. Hal ini menunjukkan bahwa data berhasil ditulis dari ApsaraMQ for Kafka ke tabel partisi di MaxCompute.