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
Anda telah mengaktifkan MNS dan memberikan izin yang diperlukan. Untuk akses jaringan internal, antrian MNS harus berada di wilayah yang sama dengan ruang kerja Flink Anda.
Jika Anda mengakses MNS melalui internet publik, Anda harus mengaktifkan akses publik untuk ruang kerja Flink Anda dan menambahkan alamat IP publik ruang kerja Flink ke daftar izin (allowlist) antrian MNS. Untuk informasi selengkapnya, lihat Kontrol Akses MNS.
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
RAWatauOSS.RAW(bawaan): Isi pesan diperlakukan sebagai JSON standar dan diurai oleh deserializerformatyang dikonfigurasi.OSS: Secara otomatis mengurai data JSON dari notifikasi event OSS. (Diperlukan langganan MNS ke event OSS.) Konektor secara otomatis mengekstraksi elemen pertama dari arrayevents, lalu meratakan dan memetakan bidang-bidang bersarang ke struktur tabel.CatatanSaat
messageTypediatur keOSS,formatharus diatur kejson.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.
CatatanKonektor 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:%';