Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:ApsaraMQ for RocketMQ

Última atualização: Jun 27, 2026

Este tópico apresenta o conector do ApsaraMQ for RocketMQ.

Importante

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

Métricas

  • Tabela source

    • numRecordsIn

    • numRecordsInPerSecond

    • numBytesIn

    • numBytesInPerSecond

    • currentEmitEventTimeLag

    • currentFetchEventTimeLag

    • sourceIdleTime

  • Tabela sink

    • numRecordsOut

    • numRecordsOutPerSecond

    • numBytesOut

    • numBytesOutPerSecond

    • currentSendTime

Nota

Consulte Métricas para mais detalhes.

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

  • Para RocketMQ 4.x, o valor é mq.

  • Para RocketMQ 5.x, o valor é mq5.

endPoint

O endpoint do serviço.

String

Sim

Nenhum

O ApsaraMQ for RocketMQ fornece dois tipos de endpoints:

  • Endpoint para o serviço MQ na rede interna (rede clássica ou VPC): Na página de detalhes da instância de destino no console MQ, selecione Endpoints > TCP Protocol Client Endpoints > Internal Network Access para obter o endpoint correspondente.

  • Endpoint público do serviço MQ: Na página de detalhes da instância de destino no console MQ, selecione Endpoint > TCP Protocol > Client Endpoint > Public Access para obter o endpoint correspondente.

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.

  • Redes internas não oferecem suporte a acesso entre regiões. Por exemplo, se o seu serviço Realtime Compute for Apache Flink estiver na região China (Hangzhou) e sua instância do Alibaba Cloud Message Queue for Apache RocketMQ estiver na região China (Shanghai), a conexão falhará.

  • Para conectar-se pela rede pública, ative o acesso público para a instância. Para mais informações, consulte Conexões de rede.

topic

O nome do tópico.

String

Sim

Nenhum

Nenhum

accessId

  • Para RocketMQ 4.x: O AccessKey ID da sua conta Alibaba Cloud.

  • Para RocketMQ 5.x:

    O nome de usuário da instância do RocketMQ.

String

  • Para RocketMQ 4.x: Sim

  • Para RocketMQ 5.x: Não

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.

  • RocketMQ 5.x:

    • Você usa um endpoint público.

    • Você usa um endpoint de VPC e o acesso sem autenticação pela rede interna está desativado.

    • Esta configuração não é necessária se você usar um endpoint de VPC e o acesso à rede interna sem autenticação estiver ativado.

accessKey

  • Para RocketMQ 4.x: O AccessKey Secret da sua conta Alibaba Cloud.

  • Para RocketMQ 5.x: A senha da instância.

String

  • Para RocketMQ 4.x: Sim

  • Para RocketMQ 5.x: Não

Nenhum

tag

A tag da mensagem para assinar ou gravar.

String

Não

Nenhum

  • Quando o RocketMQ é usado como source, é possível ler mensagens com uma única tag.

  • Ao utilizar o RocketMQ como sink, especifique várias tags separadas por vírgulas (,).

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

  • Se a instância não possuir um namespace dedicado, não configure o parâmetro instanceID.

  • Caso a instância tenha um namespace dedicado, o parâmetro instanceID será obrigatório.

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 startMessageOffset, que terá precedência.

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:

  • No modo somente leitura (padrão), o delimitador é \u0001. O delimitador não fica visível neste modo.

  • No modo de edição, o delimitador é ^A.

lengthCheck

A política para verificar o número de campos em cada registro.

String

Não

NONE

Valores válidos:

  • NONE: Valor padrão.

    • Se um registro tiver mais campos do que o esquema define, os campos extras à direita serão truncados.

    • Se um registro tiver menos campos do que o esquema define, o registro será ignorado.

  • SKIP: Ignora qualquer registro cuja contagem de campos não corresponda ao esquema.

  • EXCEPTION: Lança uma exceção se a contagem de campos não corresponder ao esquema.

  • PAD: Preenche os campos da esquerda para a direita.

    • Se um registro tiver mais campos do que o esquema define, os campos extras à direita serão truncados.

    • Se um registro tiver menos campos do que o esquema define, os campos ausentes à direita serão preenchidos com valores nulos.

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 mode estiver definido como partition.

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 deliveryTimestampValue para determinar quando uma mensagem atrasada será entregue.

String

Não

Nenhum

Valores válidos:

  • fixed: modo de timestamp fixo.

  • relative: modo de atraso relativo.

  • field: modo baseado em campo.

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:

  • event_time: tempo do evento.

  • processing_time: tempo de processamento.

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 deliveryTimestampMode:

  • deliveryTimestampMode=fixed: A mensagem é atrasada até o timestamp especificado em milissegundos. Se o horário atual for posterior ao timestamp, a mensagem será entregue imediatamente.

  • deliveryTimestampMode=relative: A duração do atraso, em milissegundos, relativa à referência de tempo especificada por deliveryTimestampType.

  • deliveryTimestampMode=field: Este parâmetro é ignorado. O horário de entrega é determinado pelo valor do campo especificado por deliveryTimestampField.

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 BIGINT.

String

Não

Nenhum

Entra em vigor quando deliveryTimestampMode é field.

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,female
    Nota

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

      Para 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 keys e tags para 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

Importante

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

Nota

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 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>
Nota

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?