Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur Simple Message Queue (anciennement MNS)

Dernière mise à jour :Aug 09, 2026

Le connecteur SMQ permet aux jobs Flink de consommer en temps réel des messages depuis les files d'attente Simple Message Queue (anciennement MNS), y compris les notifications d'événements OSS poussées par MNS.

Présentation

Par exemple, dans un pipeline de traitement d'images en temps réel, le connecteur SMQ détecte les fichiers nouvellement téléchargés dans un bucket OSS, la fonction FETCH_CONTENT télécharge l'image et l'intégration du modèle de langage IA réalise une analyse multimodale.

Catégorie

Détails

Type pris en charge

Source SQL

Mode d'exécution

Stream

Format des données

ORC, Parquet, Avro, CSV, JSON et brut

Métriques de surveillance

duplicateMessages, deletedMessages, failedDeletes, deserializationErrors

Type d'API

Flink SQL

Prérequis

Avant de commencer, vérifiez que vous avez :

Limites

  • Le connecteur SMQ requiert Ververica Runtime (VVR) version 11.6.0 ou ultérieure.

  • Pas de consommation basée sur l'offset : Contrairement à Apache Kafka, le connecteur SMQ ne prend pas en charge la consommation à partir d'un offset spécifique ni la recherche d'offset.

  • Parallélisme fixe de 1 : Le connecteur SMQ garantit la sémantique exactly-once grâce à la déduplication côté connecteur et ne prend en charge qu'un parallélisme de 1.

  • Points de contrôle requis : Le connecteur s'appuie sur les points de contrôle pour acquitter et supprimer les messages. Sans points de contrôle, les messages ne sont jamais supprimés, ce qui entraîne une consommation répétée infinie.

  • Limite de taille du corps du message : Le corps d'un seul message MNS ne peut pas dépasser 64 Ko. Pour les messages plus volumineux, consultez la section Transmettre des messages de grande taille.

  • Configuration du délai de visibilité : Définissez le délai de visibilité sur une valeur supérieure à l'intervalle des points de contrôle afin d'éviter la redistribution des messages.

  • Aucun ordre strict : Les files d'attente standard MNS ne respectent pas strictement le principe Premier entré, premier sorti (FIFO). La redistribution des messages après expiration du délai peut modifier l'ordre de consommation.

Syntaxe

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'
);

Options du connecteur

Général

Option

Type

Obligatoire

Valeur par défaut

Description

connector

String

Oui

-

Le type de connecteur. Fixé à mns.

endpoint

String

Oui

-

Le point de terminaison du service MNS.

Format : http://{account-id}.mns.{region}.aliyuncs.com. Pour plus de détails, consultez la page Points de terminaison.

region

String

Oui

-

La région où réside la file d'attente SMQ.

Exemple : ap-southeast-1. Pour plus de détails, consultez la page Points de terminaison.

queueName

String

Oui

-

Le nom de la file d'attente MNS à partir de laquelle consommer.

accessKeyId

String

Oui

-

L'ID AccessKey pour l'authentification SMQ.

Utilisez un AccessKey existant ou créez une paire AccessKey.

accessKeySecret

String

Oui

-

Le secret AccessKey pour l'authentification MNS.

Utilisez un AccessKey existant ou créez une paire AccessKey.

format

String

Oui

-

Le format de désérialisation des messages.

Valeurs valides : csv, json, avro, parquet, orc, raw.

batchSize

Integer

Non

1

Le nombre maximal de messages par demande d'interrogation par lot.

Plage valide : 1–16. Des valeurs plus élevées augmentent le débit. L'API SMQ autorise jusqu'à 16 messages par appel.

pollingWaitTime

Duration

Non

10s

Le temps d'attente de l'interrogation longue lorsque aucun message n'est disponible dans la file d'attente.

Plage valide : 0s–30s. 0s désactive l'interrogation longue.

messageType

String

Non

RAW

Le type de charge utile du message. Valeurs valides :

  • RAW

  • OSS

deleteMaxRetries

Integer

Non

3

Le nombre maximal de tentatives en cas d'échec de la suppression du message.

startTimeMs

Long

Non

-1

L'heure de début de la consommation sous forme d'horodatage Unix en millisecondes. -1 désactive le filtrage par horodatage et consomme tous les messages visibles. Contrairement à Kafka, cette option filtre les messages uniquement par heure et ne prend pas en charge la relecture des messages.

Valeurs de messageType :

  • RAW : Traite le corps du message comme des données standard et le désérialise en utilisant le paramètre format configuré.

  • OSS : Analyse le JSON issu des notifications d'événements OSS. Nécessite un abonnement MNS aux événements OSS. Le connecteur extrait le premier élément du tableau events, puis aplatit et mappe les champs imbriqués vers le schéma de la table.

Remarque

Lorsque messageType est défini sur OSS, format doit être défini sur json.

Lire les notifications d'événements OSS

Lorsque 'messageType' = 'OSS', le connecteur SMQ analyse automatiquement le JSON de notification d'événement OSS et mappe les champs à une structure de table plate. Les mappages de noms de champs suivants sont pris en charge (sans distinction entre majuscules et minuscules) :

Nom du champ

Chemin JSON

Description

eventName

eventName

Le type d'événement.

eventSource

eventSource

La source de l'événement.

eventTime

eventTime

L'heure de l'événement.

eventVersion

eventVersion

La version du protocole d'événement.

region

region

La région du bucket.

ossBucketArn

oss.bucket.arn

L'identifiant unique du bucket.

ossBucketName

oss.bucket.name

Le nom du bucket.

ossBucketOwnerIdentity

oss.bucket.ownerIdentity

L'ID utilisateur du créateur du bucket.

ossObjectKey

oss.object.key

Le nom de l'objet.

ossObjectSize

oss.object.size

La taille de l'objet.

ossObjectETag

oss.object.eTag

L'ETag de l'objet, utilisé pour détecter les modifications de contenu.

ossObjectDeltaSize

oss.object.deltaSize

La variation de la taille de l'objet.

ossObjectReadFrom

oss.object.readFrom

La position de début de lecture du fichier.

ossObjectReadTo

oss.object.readTo

La position de fin de lecture du fichier.

ossOssSchemaVersion

oss.ossSchemaVersion

Le numéro de version du schéma OSS.

ossRuleId

oss.ruleId

L'ID de la règle correspondante.

requestParametersSourceIPAddress

requestParameters.sourceIPAddress

L'adresse IP source de la requête.

responseElementsRequestId

responseElements.requestId

L'ID de la requête.

userIdentityPrincipalId

userIdentity.principalId

L'ID utilisateur (UID) du demandeur.

Les paramètres de format JSON suivants contrôlent le comportement d'analyse :

  • json.fail-on-missing-field : Indique si l'opération échoue lorsqu'un champ est manquant. Valeur par défaut : false.

  • json.ignore-parse-errors : Indique si les erreurs d'analyse sont ignorées. Valeur par défaut : false.

  • json.timestamp-format.standard : La norme de format d'horodatage. Valeur par défaut : ISO-8601.

  • json.timestamp-format.pattern : Un modèle personnalisé pour le format d'horodatage.

Pour des descriptions détaillées des champs, consultez la section Notifications d'événements OSS.

Exemples

Exemple 1 : Consommer des messages JSON standard

-- Create an MNS source table.
-- Field names must match the keys in the JSON message body.
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'
);

-- Create a print sink to test output.
CREATE TEMPORARY TABLE print_sink (
  `userId` BIGINT,
  `action` STRING,
  `timestamp` TIMESTAMP(3),
  `payload` STRING
) WITH (
  'connector' = 'print'
);

-- Consume messages and write to the sink.
INSERT INTO print_sink
SELECT userId, action, timestamp, payload
FROM mns_source;

Exemple 2 : Consommer des messages de notification d'événement OSS

-- Create an MNS source table to automatically parse OSS event notification JSON.
-- For field name mappings, see the "Read OSS event notifications" section.
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'
);

-- Create a print sink to test output.
CREATE TEMPORARY TABLE print_sink (
  `eventName` STRING,
  `ossBucketName` STRING,
  `ossObjectKey` STRING,
  `ossObjectSize` BIGINT,
  `eventTime` TIMESTAMP(3)
) WITH (
  'connector' = 'print'
);

-- Filter for ObjectCreated events and output results.
INSERT INTO print_sink
SELECT eventName, ossBucketName, ossObjectKey, ossObjectSize, eventTime
FROM oss_event_source
WHERE eventName LIKE 'ObjectCreated:%';