All Products
Search
Document Center

ApsaraMQ for Kafka:Buat konektor OSS sink

Last Updated:Jun 22, 2026

Topik ini menjelaskan cara membuat konektor OSS sink. Anda dapat menggunakan konektor ini untuk mengekspor data dari topik sumber di ApsaraMQ for Kafka ke objek di OSS.

Prasyarat

Untuk informasi selengkapnya, lihat Prasyarat.

Catatan penggunaan

  1. Konektor mempartisi data berdasarkan waktu pemrosesan event, bukan waktu pembuatan event. Jika Anda menggunakan partisi berbasis waktu, data yang berada di dekat batas waktu mungkin dikirimkan ke direktori untuk jendela waktu berikutnya.

  2. Penanganan data kotor: Jika Anda mengonfigurasi ekspresi JSONPath untuk path partisi kustom atau konten file, tetapi data masuk tidak sesuai aturan tersebut, konektor akan mengirimkan data kotor ini ke direktori invalidRuleData/ di bucket sesuai kebijakan batching. Jika Anda menemukan direktori ini di bucket Anda, verifikasi ekspresi JSONPath Anda dan pastikan konsumen Anda tidak kehilangan data.

  3. Latensi end-to-end dapat berkisar antara beberapa detik hingga menit.

  4. Jika aturan JSONPath yang dikonfigurasi untuk path partisi kustom atau konten file perlu mengekstrak data dari body pesan sumber Kafka, Anda harus mengencode kontennya ke dalam format JSON di sumber.

  5. Konektor menulis data dari sumber hulu ke OSS secara real time dengan menambahkan ke objek. Oleh karena itu, dalam satu path partisi, objek terbaru yang terlihat biasanya masih dalam proses penulisan dan belum dalam kondisi akhir. Konsumsi data ini dengan hati-hati.

Penagihan

Tugas konektor dijalankan di Function Compute. Sumber daya komputasi yang dikonsumsi untuk pemrosesan dan transmisi data ditagih berdasarkan harga Function Compute. Untuk informasi selengkapnya, lihat Ikhtisar penagihan.

Langkah 1: Buat sumber daya layanan target

Buat bucket di Konsol OSS. Untuk informasi selengkapnya, lihat Buat bucket di konsol.

Contoh ini menggunakan bucket bernama oss-sink-connector-bucket.

Langkah 2: Buat dan mulai konektor OSS sink

  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.

    • Pembuatan Tugas

      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: Ini adalah opsi yang 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: Mengencode data biner sebagai objek JSON ke dalam payload menggunakan UTF-8.

        • Text: Mengencode data biner sebagai string UTF-8 ke dalam payload. Ini adalah format default.

        • Binary: Mengencode 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 untuk memanggil fungsi. Sistem mengagregasi 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 menyaring event. Untuk informasi selengkapnya, lihat event pattern.

      3. Pada langkah Transformation, konfigurasikan transformasi data untuk melakukan operasi seperti pemisahan, pemetaan, enrichment, dan perutean dinamis. Untuk informasi selengkapnya, lihat Gunakan Function Compute untuk membersihkan data pesan.

      4. Pada langkah Sink, pilih Object Storage Service sebagai Service Type, konfigurasikan parameter berikut, lalu klik Save.

        Parameter

        Deskripsi

        Contoh

        OSS bucket

        Bucket OSS yang telah Anda buat.

        Penting
        • Pastikan bucket yang ditentukan ada dan tidak dihapus selama tugas berjalan.

        • Kelas penyimpanan bucket harus Standard atau Infrequent Access (IA). Penyimpanan Archive tidak didukung.

        • Setelah Anda membuat tugas konektor OSS sink, sistem akan menghasilkan path file sistem .tmp/ di direktori root bucket. Jangan menghapus path ini atau menggunakan objek OSS di dalamnya.

        oss-sink-connector-bucket

        Storage Path

        a/b/c/a.txta/b/c/a.txt{millisecond Unix timestamp}_{8-character random string}Kunci objek OSS terdiri dari jalur dan nama. Contohnya, jika kunci objek adalah , jalurnya adalah dan namanya adalah . Anda dapat menyesuaikan jalur partisi. Nama tersebut dihasilkan secara otomatis oleh konektor dalam format: , misalnya, .

        • Jika Anda biarkan kosong atau atur ke /, tidak ada partisi yang diterapkan, dan data disimpan ke direktori root bucket.

        • Mendukung variabel waktu: {yyyy}, {MM}, {dd}, dan {HH} masing-masing merepresentasikan tahun, bulan, hari, dan jam. Variabel ini peka huruf besar/kecil.

        • Mendukung ekspresi JSONPath untuk menyesuaikan path, misalnya {$.data.topic} dan {$.data.partition}. Variabel JSONPath harus merupakan ekspresi JSONPath yang valid. Karena batasan path OSS, nilai yang diekstraksi menggunakan JSONPath harus bertipe int atau string. Nilainya hanya boleh berisi karakter UTF-8 standar dan tidak boleh mengandung spasi, .., emoji, /, atau \. Jika tidak, dapat terjadi exception saat menulis data.

        • Mendukung konstanta.

          Catatan

          Partisi membantu mengelompokkan data secara logis dan mencegah masalah performa akibat terlalu banyak objek kecil dalam satu path.

          Throughput konektor berkorelasi positif dengan jumlah partisi. Partisi rendah atau tidak ada partisi dapat mengakibatkan throughput rendah dan menyebabkan backlog data di sumber. Terlalu banyak partisi dapat menyebabkan fragmentasi data, meningkatkan operasi tulis, dan menghasilkan terlalu banyak objek kecil. Oleh karena itu, strategi partisi yang tepat sangat penting. Pertimbangkan rekomendasi berikut:

          • Sumber Kafka: Mendukung partisi berdasarkan waktu dan partisi. Jika performa tidak mencukupi, Anda dapat menambah jumlah partisi Kafka untuk secara tidak langsung meningkatkan throughput konektor. Contoh: prefix/{yyyy}/{MM}/{dd}/{HH}/{$.data.partition}/

          • Pengelompokan Bisnis: Mempartisi data berdasarkan bidang bisnis tertentu. Laju throughput kemudian ditentukan oleh jumlah nilai unik dari bidang ini. Contoh: prefixV2/{$.data.body.field}/

          Kami merekomendasikan menggunakan awalan konstanta berbeda untuk tugas berbeda agar menghindari beberapa tugas menulis ke direktori yang sama, yang dapat menyebabkan kebingungan data.

        alikafka_post-cn-9dhsaassdd****/guide-oss-sink-topic/yyyy/MM/dd/HH

        Batch aggregation object size

        Ukuran target untuk agregasi objek. Nilainya dalam MB. Nilai valid: 1 hingga 1.024.

        Catatan
        • Konektor menulis data dalam batch hingga 16 MB. Akibatnya, ukuran objek akhir mungkin melebihi nilai yang dikonfigurasi hingga 16 MB.

        • Untuk skenario lalu lintas tinggi, kami merekomendasikan mengatur ukuran agregasi objek batch menjadi ratusan megabyte (misalnya, 128 MB atau 512 MB) dan jendela waktu ke tingkat per jam (misalnya, 60 menit atau 120 menit).

        5

        Batch aggregation time window

        Jendela waktu untuk agregasi. Nilainya dalam menit. Nilai valid: 1 hingga 1.440.

        1

        File Compression

        • No Compression Required: Menghasilkan objek OSS tanpa ekstensi file.

        • GZIP: Menghasilkan objek dengan ekstensi .gz.

        • Snappy: Menghasilkan objek dengan ekstensi .snappy.

        • Zstd: Menghasilkan objek dengan ekstensi .zstd.

        Jika Anda memilih opsi kompresi, konektor akan mengumpulkan data berdasarkan ukuran sebelum kompresi. Akibatnya, ukuran objek di OSS akan lebih kecil daripada Batch aggregation object size yang dikonfigurasi. Setelah didekompresi, ukurannya akan mendekati nilai yang dikonfigurasi.

        No Compression Required

        File Content

        • Complete Data: Konektor membungkus pesan asli dalam protokol CloudEvents. Data lengkap mencakup data dengan pembungkus protokol CloudEvents. Dalam contoh berikut, field data berisi data pesan, dan field lainnya adalah metadata yang ditambahkan oleh protokol CloudEvents.

          {
            "specversion": "1.0",
            "id": "8e215af8-ca18-4249-8645-f96c1026****",
            "source": "acs:alikafka",
            "type": "alikafka:Topic:Message",
            "subject": "acs:alikafka:alikafka_pre-cn-i7m2msb9****:topic:****",
            "datacontenttype": "application/json; charset=utf-8",
            "time": "2022-06-23T02:49:51.589Z",
            "aliyunaccountid": "182572506381****",
            "data": {
              "topic": "****",
              "partition": 7,
              "offset": 25,
              "timestamp": 1655952591589,
              "headers": {
                "headers": [],
                "isReadOnly": false
              },
              "key": "keytest",
              "value": "hello kafka msg"
            }
          }
        • Data Extraction: Mengirimkan hanya bagian data yang diekstrak menggunakan ekspresi JSONPath. Misalnya, jika Anda menentukan $.data, hanya nilai field data yang dikirimkan ke OSS.

        Untuk menghemat biaya penyimpanan dan meningkatkan efisiensi, gunakan Data Extraction dengan ekspresi $.data. Ini hanya mengirimkan pesan sumber asli ke OSS, tanpa wrapper CloudEvents.

        Data Extraction

        $.data
    • Properti tugas

      Konfigurasikan kebijakan retry untuk pengiriman event yang gagal dan cara menangani error. Untuk informasi selengkapnya, lihat Retry dan antrian dead-letter.

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

  5. Pada kotak dialog Note, baca pesannya lalu klik OK.

    Tugas memerlukan waktu 30 hingga 60 detik untuk mulai berjalan. Anda dapat memantau perkembangannya di kolom Status pada halaman Tasks.

Langkah 3: Uji konektor OSS sink

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

  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.

    Pilih tab Console. Di field Message Key, masukkan oss-sink-k2. Di field Message Content, masukkan oss-sink-v2. Untuk Send to Specified Partition, pilih No.

  4. Pada halaman Tasks, klik bucket target di kolom Event Target untuk tugas Anda.

  5. Pada halaman bucket, di panel navigasi sebelah kiri, pilih Object Management > Objects.

    • direktori tmp: Ini adalah path sistem yang digunakan oleh konektor. Jangan menghapus atau menggunakan objek OSS di path ini.

    • Direktori file data: Subdirektori dibuat berdasarkan aturan path partisi tugas, dan objek data diunggah ke direktori terdalam.

    Pada contoh ini, path breadcrumb-nya adalah /alikafka_<topic>/<partition>/2023/04/18/02/. Ini menunjukkan bahwa subdirektori dibuat berdasarkan nama topik, nomor partisi, dan struktur tahun/bulan/hari/jam. Direktori terdalam berisi objek data seperti .oss_meta_file dan partition_3_of....

  6. Di kolom Actions di sebelah kanan objek, pilih 图标 > Download.

  7. Buka file yang diunduh untuk melihat konten pesan.

    {"topic":"guide-oss-sink-topic","partition":0,"offset":0,"timestamp":1681378474218,"headers":{"headers":[],"isReadOnly":false},"key":"oss-sink-k2","value":"oss-sink-v2"}
    {"topic":"guide-oss-sink-topic","partition":0,"offset":1,"timestamp":1681378491498,"headers":{"headers":[],"isReadOnly":false},"key":"oss-sink-k2","value":"oss-sink-v2"}
    {"topic":"guide-oss-sink-topic","partition":0,"offset":2,"timestamp":1681378492515,"headers":{"headers":[],"isReadOnly":false},"key":"oss-sink-k2","value":"oss-sink-v2"}

    Output berisi beberapa pesan, dengan setiap pesan diformat sebagai objek JSON pada baris baru.