Cette rubrique présente le connecteur ApsaraMQ for RocketMQ.
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 |
|
|
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 |
|
|
endPoint |
Le point de terminaison du service. |
String |
Oui |
Aucune |
ApsaraMQ for RocketMQ propose deux types de points de terminaison :
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.
|
|
topic |
Le nom du topic. |
String |
Oui |
Aucune |
Aucune |
|
accessId |
|
String |
|
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.
|
|
accessKey |
|
String |
|
Aucune |
|
|
tag |
Le tag de message auquel s'abonner ou à écrire. |
String |
Non |
Aucune |
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 |
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 |
|
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 :
|
|
lengthCheck |
La stratégie de vérification du nombre de champs dans chaque enregistrement. |
String |
Non |
NONE |
Valeurs valides :
|
|
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 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 |
String |
Non |
Aucune |
Valeurs valides :
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 :
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
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 |
String |
Non |
Aucune |
Prend effet lorsque 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,femaleRemarqueUn 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>' );RemarquePour 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
keysettagsà 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
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
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 4.x : Connecteur DataStream MQ.
ApsaraMQ for RocketMQ 5.x : Connecteur DataStream MQ.
<!--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>
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 ?