All Products
Search
Document Center

Realtime Compute for Apache Flink:Message Service (MNS)

Last Updated:Apr 11, 2026

Topik ini menjelaskan cara menggunakan konektor Message Service (MNS).

Latar Belakang

Message Service (MNS) adalah layanan messaging terdistribusi yang efisien, andal, aman, mudah digunakan, dan dapat diskalakan secara elastis. MNS menyediakan kemampuan notifikasi event untuk Object Storage Service (OSS). Dengan membuat aturan notifikasi event, Anda dapat mendorong pesan ke antrian MNS ketika event tertentu, seperti pembuatan objek, terjadi pada sumber daya OSS yang ditentukan.

Pekerjaan Flink dapat menggunakan konektor MNS untuk mengonsumsi event tersebut. Sebagai contoh, dalam skenario pemrosesan gambar real-time, Anda dapat menggunakan konektor MNS untuk mendapatkan path file baru di bucket OSS. Selanjutnya, Anda dapat memanfaatkan fungsi FETCH_CONTENT yang disediakan oleh Realtime Compute for Apache Flink untuk mengunduh konten gambar dan menggunakan integrasi model bahasa besar AI guna melakukan analisis multimodal real-time.

Kategori

Detail

Jenis yang didukung

tabel sumber

Mode eksekusi

Streaming

Format data

Orc, Parquet, Avro, CSV, JSON, dan Raw

Metrik pemantauan spesifik

duplicateMessages (jumlah pesan duplikat), deletedMessages (jumlah pesan yang dihapus), failedDeletes (jumlah kegagalan penghapusan), deserializationErrors (jumlah error deserialisasi)

Jenis API

SQL

Prasyarat

Batasan

  • Konektor MNS hanya didukung di Realtime Compute for Apache Flink VVR 11.6.0 dan versi setelahnya.

  • Berbeda dengan Kafka, konektor MNS tidak mendukung konsumsi pesan dari offset tertentu atau melakukan seek. Untuk informasi selengkapnya, lihat Message Service (MNS).

  • Paralelisme tetap sebesar 1: Konektor MNS mencapai semantik Exactly-Once melalui deduplikasi di sisi konektor dan hanya mendukung paralelisme sebesar 1.

  • Checkpoint harus diaktifkan: Konektor MNS bergantung pada checkpoint untuk mengakui dan menghapus pesan. Jika checkpoint tidak diaktifkan, pesan tidak akan dihapus, sehingga menyebabkan konsumsi berulang tanpa henti.

  • Batas ukuran isi pesan: Isi satu pesan MNS tidak boleh melebihi 64 KB. Untuk menangani pesan yang lebih besar, lihat praktik terbaik untuk Transmisi Pesan Besar dalam dokumentasi MNS.

  • Pengaturan visibility timeout: Saat membuat antrian MNS, Anda harus mengatur visibility timeout. Atur visibility timeout ke nilai yang lebih besar daripada interval checkpoint untuk menghindari pengiriman ulang pesan.

  • Konektor MNS tidak menjamin pengurutan ketat event yang dikonsumsi. Antrian standar MNS tidak bersifat First-In, First-Out (FIFO) secara ketat. Pengiriman ulang pesan setelah timeout dapat mengubah urutan konsumsi.

Sintaks

CREATE TABLE mns_source (
    data STRING
) WITH (
    'connector' = 'mns',
    'endpoint' = '${endpoint}',
    'region' = '${region}',
    'queueName' = '${queueName}',
    'accessKeyId' = '${accessKeyId}',
    'accessKeySecret' = '${accessKeySecret}',
    'format' = 'json',
    'batchSize' = '8',
    'pollingWaitTime' = '10s',
    'messageType' = 'RAW'
);

Parameter WITH

  • Umum

    Parameter

    Deskripsi

    Tipe

    Wajib

    Bawaan

    Keterangan

    connector

    Jenis konektor.

    String

    Ya

    Tidak ada

    Nilainya tetap mns.

    endpoint

    Titik akhir layanan MNS.

    String

    Ya

    Tidak ada

    URI memiliki format berikut: http://{account-id}.mns.{region}.aliyuncs.com.

    Untuk informasi selengkapnya, lihat Titik Akhir Wilayah.

    region

    Wilayah antrian MNS.

    String

    Ya

    Tidak ada

    Contoh: cn-hangzhou. Untuk daftar wilayah yang didukung, lihat Titik Akhir.

    queueName

    Nama antrian MNS.

    String

    Ya

    Tidak ada

    Nama yang ditetapkan saat membuat antrian di konsol MNS.

    accessKeyId

    ID AccessKey untuk mengakses layanan MNS.

    String

    Ya

    Tidak ada

    Gunakan AccessKey yang sudah ada atau Buat AccessKey.

    accessKeySecret

    Rahasia AccessKey untuk mengakses layanan MNS.

    String

    Ya

    Tidak ada

    format

    Format data.

    String

    Ya

    Tidak ada

    Nilai yang valid:

    • csv

    • json

    • avro

    • parquet

    • orc

    • raw

    batchSize

    Jumlah maksimum pesan yang ditarik dari antrian MNS per batch.

    Integer

    Tidak

    1

    Rentang: 1 hingga 16. Ini memengaruhi performa baca. Nilai yang lebih besar dapat meningkatkan throughput. Catatan: API layanan MNS membatasi satu panggilan maksimal 16 pesan.

    pollingWaitTime

    Waktu maksimum menunggu pesan saat polling antrian MNS.

    Duration

    Tidak

    10s

    Rentang: 0s hingga 30s. Nilai 0s berarti tidak menunggu.

    messageType

    Tipe muatan pesan MNS.

    String

    Tidak

    RAW

    Nilainya dapat berupa RAW atau OSS.

    RAW (bawaan): Isi pesan diperlakukan sebagai JSON standar dan diurai oleh deserializer format yang dikonfigurasi.

    OSS: Secara otomatis mengurai data JSON dari notifikasi event OSS. (Diperlukan langganan MNS ke event OSS.) Konektor secara otomatis mengekstraksi elemen pertama dari array events, lalu meratakan dan memetakan bidang-bidang bersarang ke struktur tabel.

    Catatan

    Saat messageType diatur ke OSS, format harus diatur ke json.

    deleteMaxRetries

    Jumlah maksimum percobaan ulang jika penghapusan pesan gagal.

    Integer

    Tidak

    3

    Tidak ada.

    startTimeMs

    Waktu mulai konsumsi, dalam bentuk timestamp Unix dalam milidetik.

    Long

    Tidak

    -1

    Nilai -1 menonaktifkan penyaringan berdasarkan timestamp, sehingga konektor mengonsumsi semua pesan yang terlihat.

    Catatan

    Konektor MNS tidak mendukung fitur seek seperti Kafka. Parameter startTimeMs hanya digunakan untuk penyaringan. Konektor akan menyaring pesan yang masuk ke antrian sebelum timestamp ini.

Baca notifikasi event OSS

Saat Anda mengatur 'messageType' = 'OSS', konektor MNS secara otomatis mengurai bidang-bidang JSON dalam notifikasi event OSS dan memetakannya ke struktur tabel datar. Pemetaan nama bidang berikut didukung (tidak peka huruf besar/kecil):

Nama bidang

JSON path

Deskripsi

eventName

eventName

Jenis event.

eventSource

eventSource

Sumber event.

eventTime

eventTime

Waktu event.

eventVersion

eventVersion

Versi protokol event.

region

region

Wilayah bucket.

ossBucketArn

oss.bucket.arn

Identifikasi unik bucket.

ossBucketName

oss.bucket.name

Nama bucket.

ossBucketOwnerIdentity

oss.bucket.ownerIdentity

ID pengguna pembuat bucket.

ossObjectKey

oss.object.key

Nama objek.

ossObjectSize

oss.object.size

Ukuran objek.

ossObjectETag

oss.object.eTag

ETag objek, digunakan untuk memeriksa perubahan konten.

ossObjectDeltaSize

oss.object.deltaSize

Perubahan ukuran objek.

ossObjectReadFrom

oss.object.readFrom

Posisi awal pembacaan file.

ossObjectReadTo

oss.object.readTo

Posisi akhir pembacaan file.

ossOssSchemaVersion

oss.ossSchemaVersion

Nomor versi skema OSS.

ossRuleId

oss.ruleId

ID aturan yang cocok.

requestParametersSourceIPAddress

requestParameters.sourceIPAddress

Alamat IP sumber permintaan.

responseElementsRequestId

responseElements.requestId

ID permintaan unik.

userIdentityPrincipalId

userIdentity.principalId

ID pengguna (UID) pemohon.

Gunakan parameter format JSON berikut untuk mengontrol perilaku penguraian:

  • json.fail-on-missing-field: Menentukan apakah proses harus gagal jika suatu bidang tidak ditemukan (bawaan: false).

  • json.ignore-parse-errors: Menentukan apakah error penguraian harus diabaikan (bawaan: false).

  • json.timestamp-format.standard: Format timestamp (bawaan: ISO-8601).

  • json.timestamp-format.pattern: Pola kustom untuk format timestamp.

Untuk deskripsi bidang yang lebih rinci, lihat Notifikasi Event OSS.

Contoh

  • Contoh 1: Mengonsumsi pesan JSON standar

    -- Buat tabel sumber MNS.
    -- Nama bidang harus sesuai dengan kunci dalam isi pesan JSON.
    CREATE TEMPORARY TABLE mns_source (
      `userId` BIGINT,
      `action` STRING,
      `timestamp` TIMESTAMP(3),
      `payload` STRING
    ) WITH (
      'connector' = 'mns',
      'endpoint' = 'http://your-account-id.mns.cn-hangzhou.aliyuncs.com',
      'region' = 'cn-hangzhou',
      'queueName' = 'my-events-queue',
      'accessKeyId' = 'your-ak',
      'accessKeySecret' = 'your-sk',
      'format' = 'json'
    );
    
    -- Buat tabel sink untuk menguji output.
    CREATE TEMPORARY TABLE print_sink (
      `userId` BIGINT,
      `action` STRING,
      `timestamp` TIMESTAMP(3),
      `payload` STRING
    ) WITH (
      'connector' = 'print'
    );
    
    -- Konsumsi dan keluarkan data.
    INSERT INTO print_sink
    SELECT userId, action, timestamp, payload
    FROM mns_source;
  • Contoh 2: Mengonsumsi pesan notifikasi event OSS

    -- Buat tabel sumber MNS untuk mengurai JSON notifikasi event OSS secara otomatis.
    -- Untuk definisi nama bidang, lihat bagian "Baca notifikasi event OSS".
    CREATE TEMPORARY TABLE oss_event_source (
      `eventName` STRING,
      `eventSource` STRING,
      `eventTime` TIMESTAMP(3),
      `region` STRING,
      `ossBucketName` STRING,
      `ossObjectKey` STRING,
      `ossObjectSize` BIGINT,
      `responseElementsRequestId` STRING,
      `userIdentityPrincipalId` STRING
    ) WITH (
      'connector' = 'mns',
      'endpoint' = 'http://123456789.mns.cn-hangzhou.aliyuncs.com',
      'region' = 'cn-hangzhou',
      'queueName' = 'oss-events-queue',
      'accessKeyId' = '${secret_values.ak_id}',
      'accessKeySecret' = '${secret_values.ak_secret}',
      'format' = 'json',
      'messageType' = 'OSS'
    );
    
    -- Buat tabel sink untuk menguji output.
    CREATE TEMPORARY TABLE print_sink (
      `eventName` STRING,
      `ossBucketName` STRING,
      `ossObjectKey` STRING,
      `ossObjectSize` BIGINT,
      `eventTime` TIMESTAMP(3)
    ) WITH (
      'connector' = 'print'
    );
    
    -- Filter untuk event ObjectCreated dan keluarkan hasilnya.
    INSERT INTO print_sink
    SELECT eventName, ossBucketName, ossObjectKey, ossObjectSize, eventTime
    FROM oss_event_source
    WHERE eventName LIKE 'ObjectCreated:%';