Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:ApsaraMQ for RocketMQ

Dernière mise à jour :Aug 09, 2026

Cette rubrique présente le connecteur ApsaraMQ for RocketMQ.

Important

Les instances ApsaraMQ for RocketMQ 4.x Standard Edition partagent une limite supérieure élastique de 5 000 appels API par seconde. Le dépassement de cette limite lors de la connexion d'une instance à Realtime Compute for Apache Flink déclenche une limitation de débit (throttling), susceptible de déstabiliser vos jobs Flink. Si vous utilisez ou envisagez d'utiliser une instance RocketMQ Standard Edition pour l'intégrer à Flink, évaluez attentivement l'impact potentiel. Privilégiez si possible un autre middleware de messagerie, tel que Kafka, Simple Log Service (SLS) ou DataHub. Si vous devez absolument utiliser une instance ApsaraMQ for RocketMQ 4.x Standard Edition pour des volumes de messages élevés, soumettez un ticket afin de demander un quota de limitation plus élevé.

Contexte

ApsaraMQ for RocketMQ est un middleware de messagerie distribué développé par Alibaba Cloud sur la base d'Apache RocketMQ. Il offre une faible latence, une haute concurrence, une haute disponibilité et une grande fiabilité. ApsaraMQ for RocketMQ permet le découplage asynchrone et l'écrêtage des pics de charge pour les systèmes d'applications distribuées. Il propose également des fonctionnalités dédiées aux applications Internet, telles que l'accumulation massive de messages, un débit élevé et des nouvelles tentatives fiables.

Le tableau suivant décrit le connecteur ApsaraMQ for RocketMQ.

Élément

Description

Type pris en charge

Table source et table sink

Mode d'exécution

mode streaming uniquement

Format de données

Formats CSV et binaire

Métriques spécifiques au connecteur

Métriques

  • Table source

    • numRecordsIn

    • numRecordsInPerSecond

    • numBytesIn

    • numBytesInPerSecond

    • currentEmitEventTimeLag

    • currentFetchEventTimeLag

    • sourceIdleTime

  • Table sink

    • numRecordsOut

    • numRecordsOutPerSecond

    • numBytesOut

    • numBytesOutPerSecond

    • currentSendTime

Remarque

Consultez la section Métriques pour plus de détails.

Type d'API

API DataStream (pour RocketMQ 4.x uniquement) et API SQL

Prise en charge des mises à jour ou suppressions de données dans une table sink

Seule l'insertion de données dans une table sink est prise en charge. Les mises à jour et les suppressions ne sont pas prises en charge.

Fonctionnalités

Les tables source et sink ApsaraMQ for RocketMQ prennent en charge les champs de métadonnées suivants.

  • Champs pour une table source

    Champ

    Type

    Description

    topic

    VARCHAR METADATA VIRTUAL

    Le topic du message.

    queue-id

    INT METADATA VIRTUAL

    L'ID de file d'attente.

    queue-offset

    BIGINT METADATA VIRTUAL

    Le décalage de consommation.

    msg-id

    VARCHAR METADATA VIRTUAL

    L'ID du message.

    store-timestamp

    TIMESTAMP(3) METADATA VIRTUAL

    L'horodatage de stockage du message.

    born-timestamp

    TIMESTAMP(3) METADATA VIRTUAL

    L'horodatage de génération du message.

    keys

    VARCHAR METADATA VIRTUAL

    Les clés du message.

    tags

    VARCHAR METADATA VIRTUAL

    Les tags du message.

  • Champs pour une table sink

    Champ

    Type

    Description

    keys

    VARCHAR METADATA

    Les clés du message.

    tags

    VARCHAR METADATA

    Les tags du message.

Prérequis

Vous avez créé une ressource Message Queue for Apache RocketMQ. Pour obtenir des instructions, consultez la section Créer des ressources.

Limites

  • ApsaraMQ for RocketMQ 5.x nécessite le moteur de calcul en temps réel Flink VVR 8.0.3 ou version ultérieure.

  • Le connecteur ApsaraMQ for RocketMQ utilise des consommateurs par extraction (pull consumers), répartissant la charge de travail sur toutes les sous-tâches.

Syntaxe

CREATE TABLE mq_source(
  x varchar,
  y varchar,
  z varchar
) WITH (
  'connector' = 'mq5',
  'topic' = '<yourTopicName>',
  'endpoint' = '<yourEndpoint>',
  'consumerGroup' = '<yourConsumerGroup>'
);

Paramètres WITH

Général

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

connector

Le type de connecteur.

String

Oui

Aucune

  • Pour RocketMQ 4.x, la valeur est mq.

  • Pour RocketMQ 5.x, la valeur est mq5.

endPoint

Le point de terminaison du service.

String

Oui

Aucune

ApsaraMQ for RocketMQ propose deux types de points de terminaison :

  • Point de terminaison pour le service MQ sur le réseau interne (réseau classique ou VPC) : sur la page de détails de l'instance cible dans la console MQ, sélectionnez Endpoints > TCP Protocol Client Endpoints > Internal Network Access pour obtenir le point de terminaison correspondant.

  • Point de terminaison public du service MQ : sur la page de détails de l'instance cible dans la console MQ, sélectionnez Endpoint > TCP Protocol > Client Endpoint > Public Access pour obtenir le point de terminaison correspondant.

Important

Nous vous recommandons d'utiliser un point de terminaison VPC. Les connexions via le réseau public peuvent être instables en raison des changements dynamiques des politiques de sécurité réseau d'Alibaba Cloud.

  • Les réseaux internes ne prennent pas en charge l'accès interrégional. Par exemple, si votre service Realtime Compute for Apache Flink se trouve dans la région Chine (Hangzhou) et que votre instance Alibaba Cloud Message Queue for Apache RocketMQ se trouve dans la région Chine (Shanghai), la connexion échoue.

  • Pour établir une connexion via le réseau public, vous devez activer l'accès public pour l'instance. Pour plus d'informations, consultez la section Connexions réseau.

topic

Le nom du topic.

String

Oui

Aucune

Aucune

accessId

  • Pour RocketMQ 4.x : l'AccessKey ID de votre compte Alibaba Cloud.

  • Pour RocketMQ 5.x :

    Le nom d'utilisateur de l'instance RocketMQ.

String

  • Pour RocketMQ 4.x : Oui

  • Pour RocketMQ 5.x : Non

Aucune

Important

Pour éviter d'exposer votre paire AccessKey, nous vous recommandons d'utiliser une variable de projet pour spécifier l'AccessKey ID et l'AccessKey Secret.

  • RocketMQ 5.x :

    • Vous utilisez un point de terminaison public.

    • Vous utilisez un point de terminaison VPC et l'accès sans authentification sur le réseau interne est désactivé.

    • Ce paramètre n'est pas requis si vous utilisez un point de terminaison VPC et que l'accès sans authentification sur le réseau interne est activé.

accessKey

  • Pour RocketMQ 4.x : l'AccessKey Secret de votre compte Alibaba Cloud.

  • Pour RocketMQ 5.x : le mot de passe de l'instance.

String

  • Pour RocketMQ 4.x : Oui

  • Pour RocketMQ 5.x : Non

Aucune

tag

Le tag de message auquel s'abonner ou à écrire.

String

Non

Aucune

  • Lorsque RocketMQ est utilisé comme source, vous pouvez lire les messages avec un seul tag.

  • Lorsque RocketMQ est utilisé comme sink, vous pouvez spécifier plusieurs tags, séparés par des virgules (,).

Remarque

Lorsqu'il est utilisé comme sink, ce paramètre n'est pris en charge que pour RocketMQ 4.x. Pour RocketMQ 5.x, spécifiez le tag du message dans les champs de métadonnées du sink.

encoding

Le format d'encodage.

String

Non

UTF-8

Aucune

instanceID

L'ID de l'instance Alibaba Cloud Message Queue for Apache RocketMQ.

String

Non

Aucune

  • Si l'instance ne dispose pas d'espace de noms dédié, ne configurez pas le paramètre instanceID.

  • Si l'instance dispose d'un espace de noms dédié, le paramètre instanceID est obligatoire.

Remarque

Ce paramètre n'est pris en charge que pour RocketMQ 4.x.

Spécifique à la source

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

consumerGroup

Le nom du groupe de consommateurs.

String

Oui

Aucune

Aucune

pullIntervalMs

L'intervalle d'interrogation en millisecondes pour la source lorsqu'aucune donnée n'est disponible.

Int

Oui

Aucune

Unité : millisecondes.

Aucun mécanisme de limitation n'est disponible. Vous ne pouvez pas définir le débit de lecture des données depuis RocketMQ.

Remarque

Ce paramètre n'est pris en charge que pour RocketMQ 4.x.

timeZone

Le fuseau horaire.

String

Non

Aucune

Exemple : Asia/Shanghai.

startTimeMs

L'heure de début pour la consommation des données.

Long

Non

Aucune

Un horodatage en millisecondes.

startMessageOffset

Le décalage de message à partir duquel commencer la consommation.

Int

Non

Aucune

Si ce paramètre est spécifié, le chargement des données commence à partir du décalage indiqué par startMessageOffset, qui est prioritaire.

lineDelimiter

Le délimiteur de ligne utilisé pour analyser les enregistrements.

String

Non

\n

Aucune

fieldDelimiter

Le délimiteur de champ.

String

Non

\u0001

Le délimiteur varie selon le mode terminal :

  • En mode lecture seule (par défaut), le délimiteur est \u0001. Le délimiteur n'est pas visible dans ce mode.

  • En mode édition, le délimiteur est ^A.

lengthCheck

La stratégie de vérification du nombre de champs dans chaque enregistrement.

String

Non

NONE

Valeurs valides :

  • NONE : la valeur par défaut.

    • Si un enregistrement contient plus de champs que définis dans le schéma, les champs supplémentaires à droite sont tronqués.

    • Si un enregistrement contient moins de champs que définis dans le schéma, l'enregistrement est ignoré.

  • SKIP : ignore tout enregistrement dont le nombre de champs ne correspond pas au schéma.

  • EXCEPTION : lève une exception si le nombre de champs ne correspond pas au schéma.

  • PAD : complète les champs de gauche à droite.

    • Si un enregistrement contient plus de champs que définis dans le schéma, les champs supplémentaires à droite sont tronqués.

    • Si un enregistrement contient moins de champs que définis dans le schéma, les champs manquants à droite sont complétés par des valeurs nulles.

columnErrorDebug

Indique s'il faut activer le mode débogage pour les erreurs d'analyse des colonnes.

Boolean

Non

false

Si la valeur est true, des journaux détaillés pour les exceptions d'analyse sont imprimés.

pullBatchSize

Le nombre maximal de messages à extraire en un seul lot.

Int

Non

64

Pris en charge dans les versions VVR 8.0.7 et ultérieures.

Spécifique au sink

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

producerGroup

Le nom du groupe de producteurs.

String

Oui

Aucune

Aucune

retryTimes

Le nombre de tentatives en cas d'échec d'une opération d'écriture.

Int

Non

10

Aucune

sleepTimeMs

L'intervalle entre les tentatives, en millisecondes.

Long

Non

5000

Aucune

partitionField

Le nom du champ à utiliser comme clé de partitionnement.

String

Non

Aucune

Ce paramètre est obligatoire si le paramètre mode est défini sur partition.

Remarque

Pris en charge dans les versions VVR 8.0.5 et ultérieures.

deliveryTimestampMode

Le mode de livraison pour les messages différés. Ce paramètre fonctionne conjointement avec le paramètre deliveryTimestampValue pour déterminer quand un message différé est livré.

String

Non

Aucune

Valeurs valides :

  • fixed : mode d'horodatage fixe.

  • relative : mode de délai relatif.

  • field : mode basé sur un champ.

Remarque

Pris en charge dans les versions VVR 11.1 et ultérieures.

deliveryTimestampType

Le type de référence temporelle pour les messages différés.

String

Non

processing_time

Valeurs valides :

  • event_time : heure de l'événement.

  • processing_time : heure de traitement.

Remarque

Pris en charge dans les versions VVR 11.1 et ultérieures.

deliveryTimestampValue

L'heure de livraison d'un message différé.

Long

Non

Aucune

La signification de ce paramètre dépend de la valeur de deliveryTimestampMode :

  • deliveryTimestampMode=fixed : le message est retardé jusqu'à l'horodatage spécifié en millisecondes. Si l'heure actuelle est postérieure à l'horodatage, le message est livré immédiatement.

  • deliveryTimestampMode=relative : la durée du délai, en millisecondes, par rapport à la référence temporelle spécifiée par deliveryTimestampType.

  • deliveryTimestampMode=field : ce paramètre est ignoré. L'heure de livraison est déterminée par la valeur du champ spécifié par deliveryTimestampField.

Remarque

Pris en charge dans les versions VVR 11.1 et ultérieures.

deliveryTimestampField

Spécifie le champ utilisé comme heure de livraison pour les messages différés. Le type de données doit être BIGINT.

String

Non

Aucune

Prend effet lorsque deliveryTimestampMode est défini sur field.

Remarque

Pris en charge dans les versions VVR 11.1 et ultérieures.

Mappage de types

Type Flink

Type RocketMQ

BOOLEAN

STRING

VARBINARY

VARCHAR

TINYINT

INTEGER

BIGINT

FLOAT

DOUBLE

DECIMAL

Exemples

Exemples de table source

  • Format CSV

    Supposons qu'un message contienne les enregistrements de données suivants au format CSV.

    1,name,male 
    2,name,female
    Remarque

    Un message Message Queue for Apache RocketMQ peut contenir zéro ou plusieurs enregistrements de données, séparés par \n.

    Utilisez l'instruction DDL suivante dans votre job Flink pour déclarer une table source Message Queue for Apache RocketMQ.

    • RocketMQ 5.x

    CREATE TABLE mq_source(
      id varchar,
      name varchar,
      gender varchar,
      topic varchar metadata virtual
    ) WITH (
      'connector' = 'mq5',
      'topic' = 'mq-test',
      'endpoint' = '<yourEndpoint>',
      'consumerGroup' = 'mq-group',
      'fieldDelimiter' = ','
    );
    • RocketMQ 4.x

    CREATE TABLE mq_source(
      id varchar,
      name varchar,
      gender varchar,
      topic varchar metadata virtual
    ) WITH (
      'connector' = 'mq',
      'topic' = 'mq-test',
      'endpoint' = '<yourEndpoint>',
      'pullIntervalMs' = '1000',
      'accessId' = '${secret_values.ak_id}',
      'accessKey' = '${secret_values.ak_secret}',
      'consumerGroup' = 'mq-group',
      'fieldDelimiter' = ','
    );
  • Format binaire

    • RocketMQ 5.x

      CREATE TEMPORARY TABLE source_table (
        mess varbinary
      ) WITH (
        'connector' = 'mq5',
        'endpoint' = '<yourEndpoint>',
        'topic' = 'mq-test',
        'consumerGroup' = 'mq-group'
      );
      
      CREATE TEMPORARY TABLE out_table (
        commodity varchar
      ) WITH (
        'connector' = 'print'
      );
      
      INSERT INTO out_table
      select 
        cast(mess as varchar)
      FROM source_table;
    • RocketMQ 4.x

      CREATE TEMPORARY TABLE source_table (
        mess varbinary
      ) WITH (
        'connector' = 'mq',
        'endpoint' = '<yourEndpoint>',
        'pullIntervalMs' = '500',
        'accessId' = '${secret_values.ak_id}',
        'accessKey' = '${secret_values.ak_secret}',
        'topic' = 'mq-test',
        'consumerGroup' = 'mq-group'
      );
      
      CREATE TEMPORARY TABLE out_table (
        commodity varchar
      ) WITH (
        'connector' = 'print'
      );
      
      INSERT INTO out_table
      select 
        cast(mess as varchar)
      FROM source_table;

Exemples de table sink

  • Créer une table sink

    • RocketMQ 5.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR
      ) WITH (
        'connector'='mq5',
        'endpoint'='<yourEndpoint>',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );
    • RocketMQ 4.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR
      ) WITH (
        'connector'='mq',
        'endpoint'='<yourEndpoint>',
        'accessId'='${secret_values.ak_id}',
        'accessKey'='${secret_values.ak_secret}',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );
      Remarque

      Pour les messages RocketMQ au format binaire, l'instruction DDL doit définir un seul champ avec le type de données VARBINARY.

  • Créer une table sink qui mappe les champs de métadonnées keys et tags à la clé et au tag du message

    • RocketMQ 5.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR,
        keys VARCHAR METADATA,
        tags VARCHAR METADATA
      ) WITH (
        'connector'='mq5',
        'endpoint'='<yourEndpoint>',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );
    • RocketMQ 4.x

      CREATE TABLE mq_sink (
        id INTEGER,
        len BIGINT,
        content VARCHAR,
        keys VARCHAR METADATA,
        tags VARCHAR METADATA
      ) WITH (
        'connector'='mq',
        'endpoint'='<yourEndpoint>',
        'accessId'='${secret_values.ak_id}',
        'accessKey'='${secret_values.ak_secret}',
        'topic'='<yourTopicName>',
        'producerGroup'='<yourGroupName>'
      );

API DataStream

Important

Pour lire et écrire des données à l'aide de l'API DataStream, utilisez le connecteur DataStream correspondant pour vous connecter à Realtime Compute for Apache Flink. Pour plus d'informations sur la configuration d'un connecteur DataStream, consultez la section Intégrer des connecteurs DataStream.

VVR fournit MetaQSource pour lire depuis RocketMQ et MetaQOutputFormat, une implémentation de OutputFormat, pour écrire dans RocketMQ. Les exemples suivants illustrent la lecture et l'écriture dans RocketMQ :

RocketMQ 5.x

Remarque

Dans ApsaraMQ for RocketMQ 5.x, la paire de clés d'accès représente le nom d'utilisateur et le mot de passe de l'instance. Vous n'avez pas besoin de configurer cette paire si vous accédez à l'instance via un réseau interne et que l'authentification ACL est désactivée.

import com.alibaba.ververica.connectors.common.sink.OutputFormatSinkFunction;
import com.alibaba.ververica.connectors.mq5.shaded.org.apache.rocketmq.common.message.MessageExt;
import com.alibaba.ververica.connectors.mq5.sink.RocketMQOutputFormat;
import com.alibaba.ververica.connectors.mq5.source.RocketMQSource;
import com.alibaba.ververica.connectors.mq5.source.reader.deserializer.RocketMQRecordDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;
import java.util.Collections;
import java.util.List;
/**
 * A demo that shows how to consume, convert, and then produce messages to ApsaraMQ for RocketMQ.
 */
public class RocketMQ5DataStreamDemo {

    public static final String ENDPOINT = "<yourEndpoint>";
    public static final String ACCESS_ID = "<accessID>";
    public static final String ACCESS_KEY = "<accessKey>";
    public static final String SOURCE_TOPIC = "<sourceTopicName>";
    public static final String CONSUMER_GROUP = "<consumerGroup>";
    public static final String SINK_TOPIC = "<sinkTopicName>";
    public static final String PRODUCER_GROUP = "<producerGroup>";

    public static void main(String[] args) throws Exception {
        // Set up the streaming execution environment
        Configuration conf = new Configuration();

        // The following two configurations are for local debugging only. Delete them before you package the job and upload it to Realtime Compute for Apache Flink.
        conf.setString("pipeline.classpaths", "file://" + "The absolute path of the uber JAR");
        conf.setString(
                "classloader.parent-first-patterns.additional",
                "com.alibaba.ververica.connectors.mq5.source.reader.deserializer.RocketMQRecordDeserializationSchema;com.alibaba.ververica.connectors.mq5.shaded.");
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);

        final DataStreamSource<String> ds =
                env.fromSource(
                        RocketMQSource.<String>builder()
                                .setEndpoint(ENDPOINT)
                                .setAccessId(ACCESS_ID)
                                .setAccessKey(ACCESS_KEY)
                                .setTopic(SOURCE_TOPIC)
                                .setConsumerGroup(CONSUMER_GROUP)
                                .setDeserializationSchema(new MyDeserializer())
                                .setStartOffset(1)
                                .build(),
                        WatermarkStrategy.noWatermarks(),
                        "source");

        ds.map(new ToMessage())
                .addSink(
                        new OutputFormatSinkFunction<>(
                                new RocketMQOutputFormat.Builder()
                                        .setEndpoint(ENDPOINT)
                                        .setAccessId(ACCESS_ID)
                                        .setAccessKey(ACCESS_KEY)
                                        .setTopicName(SINK_TOPIC)
                                        .setProducerGroup(PRODUCER_GROUP)
                                        .build()));

        env.execute();
    }

    private static class MyDeserializer implements RocketMQRecordDeserializationSchema<String> {
        @Override
        public void deserialize(List<MessageExt> record, Collector<String> out) {
            for (MessageExt messageExt : record) {
                out.collect(new String(messageExt.getBody()));
            }
        }

        @Override
        public TypeInformation<String> getProducedType() {
            return Types.STRING;
        }
    }

    private static class ToMessage implements MapFunction<String, List<MessageExt>> {

        public ToMessage() {
        }

        @Override
        public List<MessageExt> map(String s) {
            final MessageExt message = new MessageExt();
            message.setBody(s.getBytes());
            message.setWaitStoreMsgOK(true);
            return Collections.singletonList(message);
        }
    }
}

RocketMQ 4.x

import com.alibaba.ververica.connector.mq.shaded.com.alibaba.rocketmq.common.message.MessageExt;
import com.alibaba.ververica.connectors.common.sink.OutputFormatSinkFunction;
import com.alibaba.ververica.connectors.metaq.sink.MetaQOutputFormat;
import com.alibaba.ververica.connectors.metaq.source.MetaQSource;
import com.alibaba.ververica.connectors.metaq.source.reader.deserializer.MetaQRecordDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.connector.source.Boundedness;
import org.apache.flink.api.java.typeutils.ListTypeInfo;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;
import java.io.IOException;
import java.util.List;
import java.util.Properties;
import static com.alibaba.ververica.connector.mq.shaded.com.taobao.metaq.client.ExternConst.*;
/**
 * A demo that shows how to consume, convert, and then produce messages to ApsaraMQ for RocketMQ.
 */
public class RocketMQDataStreamDemo {

    public static final String ENDPOINT = "<yourEndpoint>";
    public static final String ACCESS_ID = "<accessID>";
    public static final String ACCESS_KEY = "<accessKey>";
    public static final String INSTANCE_ID = "<instanceID>";
    public static final String SOURCE_TOPIC = "<sourceTopicName>";
    public static final String CONSUMER_GROUP = "<consumerGroup>";
    public static final String SINK_TOPIC = "<sinkTopicName>";
    public static final String PRODUCER_GROUP = "<producerGroup>";

    public static void main(String[] args) throws Exception {
        // Set up the streaming execution environment
        Configuration conf = new Configuration();

        // The following two configurations are for local debugging only. Delete them before you package the job and upload it to Realtime Compute for Apache Flink.
        conf.setString("pipeline.classpaths", "file://" + "The absolute path of the uber JAR");
        conf.setString("classloader.parent-first-patterns.additional",
                "com.alibaba.ververica.connectors.metaq.source.reader.deserializer.MetaQRecordDeserializationSchema;com.alibaba.ververica.connector.mq.shaded.");
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(conf);

        // Create and add the ApsaraMQ for RocketMQ source.
        env.fromSource(createRocketMQSource(), WatermarkStrategy.noWatermarks(), "source")
                // Convert the message body to uppercase.
                .map(RocketMQDataStreamDemo2::convertMessages)
                // Create and add the ApsaraMQ for RocketMQ sink.
                .addSink(new OutputFormatSinkFunction<>(createRocketMQOutputFormat()))
                .name(RocketMQDataStreamDemo2.class.getSimpleName());
        // Compile and submit the job.
        env.execute("RocketMQ connector end-to-end DataStream demo");
    }

    private static MetaQSource<MessageExt> createRocketMQSource() {
        Properties mqProperties = createMQProperties();

        return new MetaQSource<>(SOURCE_TOPIC,
                CONSUMER_GROUP,
                null, // always null
                null, // tag of the messages to consume
                Long.MAX_VALUE, // stop timestamp in milliseconds
                -1, // start timestamp in milliseconds. Set to -1 to disable starting from an offset.
                0, // start offset
                300_000, // partition discovery interval
                mqProperties,
                Boundedness.CONTINUOUS_UNBOUNDED,
                new MyDeserializationSchema());
    }

    private static MetaQOutputFormat createRocketMQOutputFormat() {
        return new MetaQOutputFormat.Builder()
                .setTopicName(SINK_TOPIC)
                .setProducerGroup(PRODUCER_GROUP)
                .setMqProperties(createMQProperties())
                .build();
    }

    private static Properties createMQProperties() {
        Properties properties = new Properties();
        properties.put(PROPERTY_ONS_CHANNEL, "ALIYUN");
        properties.put(NAMESRV_ADDR, ENDPOINT);
        properties.put(PROPERTY_ACCESSKEY, ACCESS_ID);
        properties.put(PROPERTY_SECRETKEY, ACCESS_KEY);
        properties.put(PROPERTY_ROCKET_AUTH_ENABLED, true);
        properties.put(PROPERTY_INSTANCE_ID, INSTANCE_ID);
        return properties;
    }

    private static List<MessageExt> convertMessages(MessageExt messages) {
        return Collections.singletonList(messages);
    }

    public static class MyDeserializationSchema implements MetaQRecordDeserializationSchema<MessageExt> {
        @Override
        public void deserialize(List<MessageExt> list, Collector<MessageExt> collector) {
            for (MessageExt messageExt : list) {
                collector.collect(messageExt);
            }
        }

        @Override
        public TypeInformation<MessageExt> getProducedType() {
            return TypeInformation.of(MessageExt.class);
        }
    }
}
    }
}

XML

<!--ApsaraMQ for RocketMQ 5.x-->
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mq5</artifactId>
    <version>${vvr-version}</version>
    <scope>provided</scope>
</dependency>

<!--ApsaraMQ for RocketMQ 4.x-->
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mq</artifactId>
    <version>${vvr-version}</version>
</dependency>
Remarque

Pour plus d'informations sur la configuration du point de terminaison pour ApsaraMQ for RocketMQ, consultez la section Annonce sur les paramètres des points de terminaison TCP internes.

FAQ

Comment RocketMQ détecte-t-il les modifications du nombre de partitions lors du scaling des topics ?