Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Mensagens normais

Última atualização: Jun 27, 2026

As mensagens normais são o tipo padrão no ApsaraMQ for RocketMQ. Diferentemente das mensagens ordenadas, agendadas, atrasadas e transacionais, não possuem recursos especiais. Esse tipo suporta comunicação assíncrona e desacoplada entre produtores e consumidores. Utilize mensagens normais quando sua aplicação exigir entrega confiável, mas dispensar sequenciamento rigoroso ou entrega programada.

Quando usar mensagens normais

Mensagens normais adequam-se a qualquer cenário que priorize a confiabilidade da transmissão em vez da ordem de processamento. Exemplos comuns incluem desacoplamento de microsserviços, integração de dados e arquiteturas orientadas a eventos:

  • Desacoplamento de microsserviços: Um serviço upstream encapsula a criação do pedido e o pagamento em uma mensagem normal independente e a envia ao broker do ApsaraMQ for RocketMQ. Os serviços downstream assinam independentemente e processam a mensagem conforme sua própria lógica. Cada mensagem é autossuficiente, sem dependências cruzadas.

    Asynchronous decoupling between an order system and downstream services

  • Integração de dados: Um componente de instrumentação coleta logs de operação de aplicações frontend e os encaminha para o ApsaraMQ for RocketMQ. Cada mensagem representa um registro discreto de log. O broker armazena e entrega esses registros aos sistemas de armazenamento downstream sem transformá-los. As aplicações backend processam as tarefas subsequentes.

    Log data pipeline from frontend applications to downstream storage

Ciclo de vida da mensagem

Uma mensagem normal passa por cinco estágios entre a criação e a exclusão:

Lifecycle of a normal message

Estágio

Descrição

Inicializado

O produtor constrói a mensagem e a prepara para envio.

Pronto

O broker aceita a mensagem, tornando-a visível e disponível para os consumidores.

Em trânsito

Um consumidor recupera a mensagem e inicia o processamento. Se o broker não receber confirmação dentro de um período específico, ele tentará entregar novamente.

Confirmado

O consumidor registra o resultado do consumo. O broker marca a mensagem como consumida, mas não a exclui imediatamente.

Excluído

A remoção ocorre quando o período de retenção expira ou o espaço de armazenamento diminui. A exclusão segue uma estratégia rotativa que remove primeiro as mensagens mais antigas do arquivo físico. Para obter detalhes, consulte Armazenamento e limpeza de mensagens.

Por padrão, o ApsaraMQ for RocketMQ retém todas as mensagens mesmo após a confirmação. Os consumidores podem reconsumir qualquer mensagem ainda não excluída.

Exemplos de código (Java)

Os exemplos a seguir usam o SDK do Apache RocketMQ 5.x para enviar e consumir mensagens normais. Cada exemplo configura uma conexão de cliente e executa uma única operação. Substitua os placeholders pelos valores reais antes de executar o código.

Para consultar a referência completa do SDK e exemplos em outras linguagens, veja SDKs do Apache RocketMQ 5.x.

Antes de começar

Prepare os seguintes recursos no console do ApsaraMQ for RocketMQ:

Recurso

Onde encontrar

Endpoint

Página Instance Details > aba Endpoints. Use o endpoint de VPC para acessar instâncias do Elastic Compute Service (ECS) na mesma virtual private cloud (VPC). Utilize o endpoint público para acesso pela internet (requer ativação do recurso de acesso à internet).

Topic

Crie um tópico com MessageType definido como Normal. Mensagens normais só podem ser enviadas para tópicos desse tipo.

Consumer group

Necessário para os exemplos de consumidor. Crie um grupo no console.

Credenciais (apenas endpoint público)

Nome de usuário e senha obtidos na página Access Control > aba Intelligent Authentication. Dispensáveis para acesso via VPC, pois o broker autentica automaticamente. Também não são obrigatórias para instâncias serverless acessadas por VPC com o recurso de autenticação livre em VPCs ativado.

Referência de placeholders

Placeholder

Descrição

Exemplo

<your-endpoint>

Endpoint da instância

rmq-cn-xxx.cn-hangzhou.rmq.aliyuncs.com:8080

<your-topic>

Nome do tópico

normal-topic-01

<your-consumer-group>

Nome do grupo de consumidores

GID_normal_consumer

<your-username>

Nome de usuário da instância (apenas endpoint público)

--

<your-password>

Senha da instância (apenas endpoint público)

--

Código de exemplo

Enviar uma mensagem normal

package doc;

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.apis.producer.Producer;
import org.apache.rocketmq.client.apis.producer.SendReceipt;

public class ProducerExample {
    public static void main(String[] args) throws ClientException {
        // Specify the instance endpoint.
        String endpoints = "<your-endpoint>";
        String topic = "<your-topic>";

        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // Uncomment the following line for public endpoint access:
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        ClientConfiguration configuration = builder.build();

        Producer producer = provider.newProducerBuilder()
                .setTopics(topic)
                .setClientConfiguration(configuration)
                .build();

        // Build a normal message with a key and tag for filtering and tracing.
        Message message = provider.newMessageBuilder()
                .setTopic(topic)
                .setKeys("messageKey")
                .setTag("messageTag")
                .setBody("messageBody".getBytes())
                .build();

        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println("Send success. Message ID: " + sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

Saída esperada:

Send success. Message ID: 01BE0A3600F5762D0455025C2D1A0000

Consumir com um push consumer

Um push consumer recebe mensagens por meio de um callback MessageListener. O broker envia as mensagens automaticamente para o consumidor.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.ConsumeResult;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.PushConsumer;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.util.Collections;

public class PushConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(PushConsumerExample.class);

    public static void main(String[] args) throws ClientException, IOException, InterruptedException {
        String endpoints = "<your-endpoint>";
        String topic = "<your-topic>";
        String consumerGroup = "<your-consumer-group>";

        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // Uncomment the following line for public endpoint access:
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        ClientConfiguration clientConfiguration = builder.build();

        // Subscribe to all messages in the topic (tag = "*").
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);

        PushConsumer pushConsumer = provider.newPushConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .setMessageListener(messageView -> {
                    System.out.println("Received message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();

        // Keep the process running to continue receiving messages.
        Thread.sleep(Long.MAX_VALUE);

        // Close the consumer when it is no longer needed.
        // pushConsumer.close();
    }
}

Saída esperada:

Received message: MessageView{messageId=01BE0A3600F5762D0455025C2D1A0000, topic=<your-topic>, ...}

Consumir com um simple consumer

Um simple consumer busca mensagens explicitamente via receive() e confirma cada uma com ack(). Esse modo oferece controle total sobre o ritmo de consumo.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.SimpleConsumer;
import org.apache.rocketmq.client.apis.message.MessageId;
import org.apache.rocketmq.client.apis.message.MessageView;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.time.Duration;
import java.util.Collections;
import java.util.List;

public class SimpleConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(SimpleConsumerExample.class);

    public static void main(String[] args) throws ClientException, IOException {
        String endpoints = "<your-endpoint>";
        String topic = "<your-topic>";
        String consumerGroup = "<your-consumer-group>";

        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // Uncomment the following line for public endpoint access:
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        ClientConfiguration clientConfiguration = builder.build();

        Duration awaitDuration = Duration.ofSeconds(10);
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);

        SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setAwaitDuration(awaitDuration)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .build();

        int maxMessageNum = 16;
        Duration invisibleDuration = Duration.ofSeconds(10);

        // Pull and acknowledge messages in a loop.
        // For real-time consumption, use multiple threads to pull concurrently.
        while (true) {
            final List<MessageView> messages = consumer.receive(maxMessageNum, invisibleDuration);
            messages.forEach(messageView -> {
                System.out.println("Received message: " + messageView);
            });
            for (MessageView message : messages) {
                final MessageId messageId = message.getMessageId();
                try {
                    consumer.ack(message);
                    System.out.println("Acked message: " + messageId);
                } catch (Throwable t) {
                    t.printStackTrace();
                }
            }
        }

        // Close the consumer when it is no longer needed.
        // consumer.close();
    }
}

Saída esperada:

Received message: MessageView{messageId=01BE0A3600F5762D0455025C2D1A0000, topic=<your-topic>, ...}
Acked message: 01BE0A3600F5762D0455025C2D1A0000

Melhores práticas

Defina uma chave de mensagem globalmente única

Atribua um identificador de negócio — como ID de pedido ou de usuário — como chave da mensagem. Isso permite consultar mensagens individuais e seus rastreamentos no console do ApsaraMQ for RocketMQ.

Message message = provider.newMessageBuilder()
        .setTopic(topic)
        .setKeys("order-12345")   // Use a unique business identifier.
        .setTag("payment")
        .setBody(payload)
        .build();

Escolha o modo de consumidor adequado

Modo

Funcionamento

Indicado para

Push consumer

O broker envia mensagens para um callback MessageListener.

Processamento orientado a eventos e de baixa latência, no qual o consumidor lida com cada mensagem assim que ela chega.

Simple consumer

A aplicação chama receive() para buscar um lote e depois executa ack() em cada mensagem após o processamento.

Cargas de trabalho que exigem controle manual de fluxo ou processamento em lotes.