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 :
Simple Message Queue (anciennement MNS) : Activation et autorisations. Pour un accès via le réseau interne, la file d'attente SMQ doit se trouver dans la même région que votre espace de travail Flink.
(Pour un accès Internet public) Accès public activé pour votre espace de travail Flink et ajoutez ses adresses IP publiques à la liste d'autorisation de la file d'attente MNS. Pour plus de détails, consultez la section Contrôle d'accès.
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é à |
|
endpoint |
String |
Oui |
- |
Le point de terminaison du service MNS. Format : |
|
region |
String |
Oui |
- |
La région où réside la file d'attente SMQ. Exemple : |
|
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 : |
|
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. |
|
messageType |
String |
Non |
|
Le type de charge utile du message. Valeurs valides :
|
|
deleteMaxRetries |
Integer |
Non |
|
Le nombre maximal de tentatives en cas d'échec de la suppression du message. |
|
startTimeMs |
Long |
Non |
|
L'heure de début de la consommation sous forme d'horodatage Unix en millisecondes. |
Valeurs de messageType :
RAW: Traite le corps du message comme des données standard et le désérialise en utilisant le paramètreformatconfiguré.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 tableauevents, puis aplatit et mappe les champs imbriqués vers le schéma de la table.
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:%';