Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Send and receive normal messages

Última atualização: Jun 27, 2026

Este tópico apresenta exemplos em Java para enviar e receber mensagens normais pelo SDK de cliente HTTP do ApsaraMQ for RocketMQ. As mensagens normais são o tipo padrão: não possuem semântica especial de entrega, diferentemente das mensagens agendadas, atrasadas, ordenadas e transacionais, e chegam aos consumidores o mais rápido possível.

Pré-requisitos

Antes de começar, verifique se você já:

  • Instalou o SDK Java. Para mais informações, consulte Preparar o ambiente.

  • Crie os recursos necessários no console do ApsaraMQ for RocketMQ: uma instância, um tópico (definido como tipo de mensagem Normal) e um grupo de consumidores. Para mais informações, consulte Criar recursos.

  • Obteve seu par de AccessKey (AccessKey ID e AccessKey secret). Para mais informações, consulte Criar um par de AccessKey.

Como funciona

Ambos os exemplos seguem o mesmo fluxo de trabalho:

  1. Crie um MQClient com o endpoint HTTP e as credenciais AccessKey.

  2. Vincule um produtor ou consumidor a uma instância e a um tópico específicos.

  3. Envie mensagens ou faça polling e feche o cliente ao terminar.

Cada tópico suporta apenas um tipo de mensagem. Um tópico criado para mensagens normais não pode enviar ou receber mensagens agendadas, atrasadas, ordenadas ou transacionais.

Enviar uma mensagem

O exemplo abaixo envia quatro mensagens normais para um tópico em loop, usando publicação síncrona. Se publishMessage retornar sem lançar exceção, a mensagem foi enviada com sucesso.

Substitua os placeholders pelos valores reais:

Placeholder

Descrição

Onde encontrar

<your-http-endpoint>

Endpoint HTTP da instância

Página Instance Details > seção HTTP Endpoint

<your-topic>

Nome do tópico

Console do ApsaraMQ for RocketMQ

<your-instance-id>

ID da instância. Defina como null ou "" se a instância não tiver namespace

Página Instance Details

import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQProducer;
import com.aliyun.mq.http.model.TopicMessage;

import java.util.Date;

public class Producer {

    public static void main(String[] args) {
        MQClient mqClient = new MQClient(
                "<your-http-endpoint>",
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
        );

        final String topic = "<your-topic>";
        final String instanceId = "<your-instance-id>";

        // Get a producer for this topic and instance.
        MQProducer producer;
        if (instanceId != null && instanceId != "") {
            producer = mqClient.getProducer(instanceId, topic);
        } else {
            producer = mqClient.getProducer(topic);
        }

        try {
            for (int i = 0; i < 4; i++) {
                TopicMessage pubMsg = new TopicMessage(
                        "hello mq!".getBytes(),  // Message body
                        "A"                       // Message tag for filtering
                );
                // Custom property
                pubMsg.getProperties().put("a", String.valueOf(i));
                // Business key for message tracing (use an order ID, user ID, or similar identifier)
                pubMsg.setMessageKey("MessageKey");

                // Publish synchronously. No exception means success.
                TopicMessage pubResultMsg = producer.publishMessage(pubMsg);

                System.out.println(new Date() + " Send mq message success. Topic is:" + topic
                        + ", msgId is: " + pubResultMsg.getMessageId()
                        + ", bodyMD5 is: " + pubResultMsg.getMessageBodyMD5());
            }
        } catch (Throwable e) {
            System.out.println(new Date() + " Send mq message failed. Topic is:" + topic);
            e.printStackTrace();
        }

        mqClient.close();
    }
}

Pontos importantes:

  • Credenciais via variáveis de ambiente: O código lê ALIBABA_CLOUD_ACCESS_KEY_ID e ALIBABA_CLOUD_ACCESS_KEY_SECRET de variáveis de ambiente. Configure essas variáveis antes de executar o código.

  • Tag da mensagem: Tags permitem filtrar mensagens dentro de um tópico. Neste exemplo, a tag é "A".

  • Chave da mensagem: Defina a chave da mensagem como um identificador de negócio, como um ID de pedido ou de usuário. Isso facilita a consulta e o rastreamento de mensagens no console do ApsaraMQ for RocketMQ.

Receber e confirmar mensagens

O exemplo a seguir consome mensagens de um tópico em um loop de long polling. Após processar cada lote, o sistema envia uma confirmação (ACK) ao broker. Caso o broker não receba o ACK antes do intervalo de nova tentativa de entrega expirar, ele reentrega a mensagem.

Substitua os placeholders pelos valores reais:

Placeholder

Descrição

Onde encontrar

<your-http-endpoint>

Endpoint HTTP da instância

Página Instance Details > seção HTTP Endpoint

<your-topic>

Nome do tópico

Console do ApsaraMQ for RocketMQ

<your-instance-id>

ID da instância. Defina como null ou "" se a instância não tiver namespace

Página Instance Details

<your-group-id>

ID do grupo de consumidores

Console do ApsaraMQ for RocketMQ

import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQConsumer;
import com.aliyun.mq.http.common.AckMessageException;
import com.aliyun.mq.http.model.Message;

import java.util.ArrayList;
import java.util.List;

public class Consumer {

    public static void main(String[] args) {
        MQClient mqClient = new MQClient(
                "<your-http-endpoint>",
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
        );

        final String topic = "<your-topic>";
        final String groupId = "<your-group-id>";
        final String instanceId = "<your-instance-id>";

        // Get a consumer for this topic, group, and instance.
        final MQConsumer consumer;
        if (instanceId != null && instanceId != "") {
            consumer = mqClient.getConsumer(instanceId, topic, groupId, null);
        } else {
            consumer = mqClient.getConsumer(topic, groupId);
        }

        // Poll for messages in a loop.
        // In production, use multiple threads for concurrent consumption.
        do {
            List<Message> messages = null;

            try {
                // Long polling: if no message is available, the request hangs
                // on the broker for up to the specified number of seconds.
                messages = consumer.consumeMessage(
                        3,  // Max messages per batch (up to 16)
                        3   // Long-polling timeout in seconds (up to 30)
                );
            } catch (Throwable e) {
                e.printStackTrace();
                try {
                    Thread.sleep(2000);
                } catch (InterruptedException e1) {
                    e1.printStackTrace();
                }
            }

            if (messages == null || messages.isEmpty()) {
                System.out.println(Thread.currentThread().getName() + ": no new message, continue!");
                continue;
            }

            // --- Process messages ---
            for (Message message : messages) {
                System.out.println("Receive message: " + message);
            }

            // --- Acknowledge messages ---
            // Each message has a unique receipt handle that expires after
            // the delivery retry interval. ACK before it expires to prevent
            // redelivery.
            {
                List<String> handles = new ArrayList<String>();
                for (Message message : messages) {
                    handles.add(message.getReceiptHandle());
                }

                try {
                    consumer.ackMessage(handles);
                } catch (Throwable e) {
                    if (e instanceof AckMessageException) {
                        AckMessageException errors = (AckMessageException) e;
                        System.out.println("Ack message fail, requestId is:" + errors.getRequestId()
                                + ", fail handles:");
                        if (errors.getErrorMessages() != null) {
                            for (String errorHandle : errors.getErrorMessages().keySet()) {
                                System.out.println("Handle:" + errorHandle
                                        + ", ErrorCode:" + errors.getErrorMessages().get(errorHandle).getErrorCode()
                                        + ", ErrorMsg:" + errors.getErrorMessages().get(errorHandle).getErrorMessage());
                            }
                        }
                        continue;
                    }
                    e.printStackTrace();
                }
            }
        } while (true);
    }
}

Pontos importantes:

  • Long polling: O consumidor mantém a conexão aberta no broker durante o tempo limite especificado (3 segundos neste exemplo). Se uma mensagem chegar nesse período, o broker responde imediatamente em vez de aguardar o próximo ciclo de polling. O tempo limite máximo é de 30 segundos.

  • Tamanho do lote: Cada chamada de consumeMessage retorna até o número especificado de mensagens (3 neste exemplo, máximo de 16).

  • Receipt handle: Sempre que uma mensagem é entregue, ela recebe um novo receipt handle com um timestamp único. Use esse handle para enviar o ACK da mensagem. Se o handle expirar antes que o ACK chegue ao broker, a mensagem será reentregue.

  • Concorrência: Este exemplo usa uma única thread. Para obter maior throughput, utilize múltiplas threads para consumo em ambientes de produção.

Observações de uso

  • Um tipo por tópico: Um tópico criado para mensagens normais não pode enviar ou receber outros tipos de mensagem (agendadas, atrasadas, ordenadas ou transacionais). Crie tópicos separados para cada tipo de mensagem.

  • Defina chaves de mensagem significativas: Utilize identificadores de negócio (como IDs de pedidos ou de usuários) como chave da mensagem. Isso agiliza a consulta e o rastreamento de mensagens no console.

  • Trate falhas de ACK adequadamente: Se ackMessage lançar uma AckMessageException, registre em log os handles com falha e os códigos de erro. O broker reentrega automaticamente as mensagens não confirmadas; portanto, evite tratar falhas de ACK como erros fatais.

  • Limpeza de recursos: Chame mqClient.close() quando o produtor ou consumidor não for mais necessário para liberar as conexões HTTP.

Próximos passos