Este tópico apresenta o conector do ApsaraMQ for RocketMQ.
As instâncias da Standard Edition do ApsaraMQ for RocketMQ 4.x compartilham um limite superior elástico de 5.000 chamadas de API por segundo. Exceder esse limite ao conectar uma instância ao Realtime Compute for Apache Flink aciona o throttling, o que pode desestabilizar seus jobs do Flink. Portanto, se você usa ou planeja usar uma instância do RocketMQ Standard Edition para integração com o Flink, avalie cuidadosamente o impacto potencial. Se possível, considere um middleware de mensagens alternativo, como Kafka, Simple Log Service (SLS) ou DataHub. Caso seja indispensável utilizar uma instância da Standard Edition do ApsaraMQ for RocketMQ 4.x para altos volumes de mensagens, abra um ticket para solicitar uma cota de throttling maior.
Informações básicas
O ApsaraMQ for RocketMQ é um middleware de mensagens distribuído desenvolvido pela Alibaba Cloud com base no Apache RocketMQ. Ele oferece baixa latência, alta concorrência, alta disponibilidade e alta confiabilidade. O ApsaraMQ for RocketMQ fornece desacoplamento assíncrono e suavização de picos para sistemas de aplicações distribuídas. Também disponibiliza recursos para aplicações de internet, como acumulação massiva de mensagens, alto throughput e retentativas confiáveis.
A tabela a seguir descreve o conector do ApsaraMQ for RocketMQ.
|
Item |
Descrição |
|
Tipo suportado |
Tabela source e tabela sink |
|
Modo de execução |
Apenas modo streaming |
|
Formato de dados |
Formatos CSV e binário |
|
Métricas específicas do conector |
|
|
Tipo de API |
DataStream API (apenas para RocketMQ 4.x) e SQL API |
|
Suporte a atualizações ou exclusões de dados em uma tabela sink |
Há suporte apenas para inserção de dados na tabela sink. Não há suporte para atualizações e exclusões. |
Recursos
As tabelas source e sink do ApsaraMQ for RocketMQ aceitam os seguintes campos de metadados.
-
Campos para tabela source
Campo
Tipo
Descrição
topic
VARCHAR METADATA VIRTUAL
O tópico da mensagem.
queue-id
INT METADATA VIRTUAL
O ID da fila.
queue-offset
BIGINT METADATA VIRTUAL
O offset de consumo.
msg-id
VARCHAR METADATA VIRTUAL
O ID da mensagem.
store-timestamp
TIMESTAMP(3) METADATA VIRTUAL
O timestamp de armazenamento da mensagem.
born-timestamp
TIMESTAMP(3) METADATA VIRTUAL
O timestamp de geração da mensagem.
keys
VARCHAR METADATA VIRTUAL
As chaves da mensagem.
tags
VARCHAR METADATA VIRTUAL
As tags da mensagem.
-
Campos para tabela sink
Campo
Tipo
Descrição
keys
VARCHAR METADATA
As chaves da mensagem.
tags
VARCHAR METADATA
As tags da mensagem.
Pré-requisitos
Crie um recurso do Message Queue for Apache RocketMQ. Para obter instruções, consulte Criar recursos.
Limitações
O ApsaraMQ for RocketMQ 5.x exige o mecanismo de computação em tempo real do Flink VVR 8.0.3 ou posterior.
O conector do ApsaraMQ for RocketMQ utiliza consumidores pull e distribui a carga de trabalho entre todas as subtarefas.
Sintaxe
CREATE TABLE mq_source(
x varchar,
y varchar,
z varchar
) WITH (
'connector' = 'mq5',
'topic' = '<yourTopicName>',
'endpoint' = '<yourEndpoint>',
'consumerGroup' = '<yourConsumerGroup>'
);
Parâmetros WITH
Geral
|
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
|
connector |
O tipo do conector. |
String |
Sim |
Nenhum |
|
|
endPoint |
O endpoint do serviço. |
String |
Sim |
Nenhum |
O ApsaraMQ for RocketMQ fornece dois tipos de endpoints:
Importante
Recomendamos o uso de um endpoint de VPC. Conexões via rede pública podem ser instáveis devido a alterações dinâmicas nas políticas de segurança de rede da Alibaba Cloud.
|
|
topic |
O nome do tópico. |
String |
Sim |
Nenhum |
Nenhum |
|
accessId |
|
String |
|
Nenhum |
Importante
Para evitar expor seu par de AccessKey, recomendamos usar uma variável de projeto para especificar o AccessKey ID e o AccessKey Secret.
|
|
accessKey |
|
String |
|
Nenhum |
|
|
tag |
A tag da mensagem para assinar ou gravar. |
String |
Não |
Nenhum |
Nota
Quando usado como sink, este parâmetro tem suporte apenas no RocketMQ 4.x. Para o RocketMQ 5.x, especifique a tag da mensagem nos campos de metadados do sink. |
|
encoding |
O formato de codificação. |
String |
Não |
UTF-8 |
Nenhum |
|
instanceID |
O ID da instância do Alibaba Cloud Message Queue for Apache RocketMQ. |
String |
Não |
Nenhum |
Nota
Este parâmetro tem suporte apenas no RocketMQ 4.x. |
Específicos da source
|
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
|
consumerGroup |
O nome do grupo de consumidores. |
String |
Sim |
Nenhum |
Nenhum |
|
pullIntervalMs |
O intervalo de polling em milissegundos para a source quando não há dados disponíveis. |
Int |
Sim |
Nenhum |
Unidade: milissegundos. Não há mecanismo de throttling disponível. Não é possível definir a taxa de leitura de dados do RocketMQ. Nota
Este parâmetro tem suporte apenas no RocketMQ 4.x. |
|
timeZone |
O fuso horário. |
String |
Não |
Nenhum |
Exemplo: Asia/Shanghai. |
|
startTimeMs |
O horário inicial para consumo de dados. |
Long |
Não |
Nenhum |
Um timestamp em milissegundos. |
|
startMessageOffset |
O offset da mensagem a partir do qual iniciar o consumo. |
Int |
Não |
Nenhum |
Se este parâmetro for especificado, o carregamento de dados começará a partir do offset definido por |
|
lineDelimiter |
O delimitador de linha usado para analisar registros. |
String |
Não |
\n |
Nenhum |
|
fieldDelimiter |
O delimitador de campo. |
String |
Não |
\u0001 |
O delimitador varia conforme o modo do terminal:
|
|
lengthCheck |
A política para verificar o número de campos em cada registro. |
String |
Não |
NONE |
Valores válidos:
|
|
columnErrorDebug |
Define se o modo de depuração para erros de análise de colunas deve ser ativado. |
Boolean |
Não |
false |
Se definido como true, logs detalhados para exceções de análise serão impressos. |
|
pullBatchSize |
O número máximo de mensagens a serem buscadas em um único lote. |
Int |
Não |
64 |
Suportado no VVR 8.0.7 e versões posteriores. |
Específicos do sink
|
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
|
producerGroup |
O nome do grupo de produtores. |
String |
Sim |
Nenhum |
Nenhum |
|
retryTimes |
O número de tentativas para repetir uma operação de gravação com falha. |
Int |
Não |
10 |
Nenhum |
|
sleepTimeMs |
O intervalo entre novas tentativas, em milissegundos. |
Long |
Não |
5000 |
Nenhum |
|
partitionField |
O nome do campo a ser usado como chave de partição. |
String |
Não |
Nenhum |
Este parâmetro é necessário se o parâmetro Nota
Suportado no VVR 8.0.5 e versões posteriores. |
|
deliveryTimestampMode |
O modo de entrega para mensagens atrasadas. Este parâmetro funciona em conjunto com o parâmetro |
String |
Não |
Nenhum |
Valores válidos:
Nota
Suportado no VVR 11.1 e versões posteriores. |
|
deliveryTimestampType |
O tipo de referência de tempo para mensagens atrasadas. |
String |
Não |
processing_time |
Valores válidos:
Nota
Suportado no VVR 11.1 e versões posteriores. |
|
deliveryTimestampValue |
O horário de entrega de uma mensagem atrasada. |
Long |
Não |
Nenhum |
O significado deste parâmetro depende do valor de
Nota
Suportado no VVR 11.1 e versões posteriores. |
|
deliveryTimestampField |
Especifica o campo usado como horário de entrega para mensagens atrasadas. O tipo de dados deve ser |
String |
Não |
Nenhum |
Entra em vigor quando Nota
Suportado no VVR 11.1 e versões posteriores. |
Mapeamento de tipos
|
Tipo Flink |
Tipo RocketMQ |
|
BOOLEAN |
STRING |
|
VARBINARY |
|
|
VARCHAR |
|
|
TINYINT |
|
|
INTEGER |
|
|
BIGINT |
|
|
FLOAT |
|
|
DOUBLE |
|
|
DECIMAL |
Exemplos
Exemplos de tabela source
-
Formato CSV
Suponha que uma mensagem contenha os seguintes registros de dados no formato CSV.
1,name,male 2,name,femaleNotaUma mensagem do Message Queue for Apache RocketMQ pode conter zero ou mais registros de dados, separados por
\n.Use a seguinte DDL no seu job do Flink para declarar uma tabela source do 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' = ',' ); -
Formato binário
-
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;
-
Exemplos de tabela sink
-
Criar uma tabela 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>' );NotaPara mensagens do RocketMQ em formato binário, a DDL deve definir um único campo com o tipo de dados VARBINARY.
-
-
Criar uma tabela sink que mapeia os campos de metadados
keysetagspara a chave e a tag da mensagem-
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>' );
-
DataStream API
Para ler e gravar dados usando a DataStream API, utilize o conector DataStream correspondente para se conectar ao Realtime Compute for Apache Flink. Para mais informações sobre como configurar um conector DataStream, consulte Integrar conectores DataStream.
O VVR fornece MetaQSource para leitura do RocketMQ e MetaQOutputFormat, uma implementação de OutputFormat, para gravação no RocketMQ. Os exemplos a seguir mostram como ler e gravar no RocketMQ:
RocketMQ 5.x
No ApsaraMQ for RocketMQ 5.x, o par de chaves de acesso representa o nome de usuário e a senha da instância. Não é necessário configurar esse par se você acessar a instância por uma rede interna e a autenticação ACL estiver desativada.
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: Conector DataStream MQ.
ApsaraMQ for RocketMQ 5.x: Conector 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>
Para mais informações sobre como configurar o endpoint do ApsaraMQ for RocketMQ, consulte Comunicado sobre as configurações de endpoints internos TCP.
Perguntas frequentes
Como o RocketMQ detecta alterações na contagem de partições durante o dimensionamento do tópico?