Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Send and receive messages with the Java SDK

Última atualização: Jun 27, 2026

Após criar recursos do ApsaraMQ for RocketMQ, como um tópico e um grupo de consumidores, integre a troca de mensagens à sua aplicação Java. Este guia demonstra como adicionar a dependência do SDK 5.x para Java, conectar-se à instância, enviar mensagens normais e consumi-las com um push consumer ou um simple consumer.

Antes de começar

Verifique se você tem:

Adicione a dependência do SDK

  1. Crie um projeto Java na sua IDE.

  2. Adicione a seguinte dependência ao arquivo pom.xml:

        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-client-java</artifactId>
            <version>5.0.7</version>
        </dependency>

Reúna os parâmetros de conexão

Obtenha os seguintes valores no console do ApsaraMQ for RocketMQ antes de escrever o código:

Parâmetro

Exemplo

Onde encontrar

endpoints

rmq-cn-xxx.{regionId}.rmq.aliyuncs.com:8080

Página Instance Details > aba Endpoints. Use o endpoint de VPC para acesso interno ou o endpoint público para acesso pela Internet. Consulte Obter o endpoint de uma instância.

topic

normal_test

Tópico criado para envio e recebimento de mensagens. Crie-o previamente. Consulte Criar um tópico.

consumerGroup

GID_test

Grupo de consumidores criado. Crie-o previamente. Consulte Criar um grupo de consumidores.

InstanceId

rmq-cn-xxx

ID da instância do ApsaraMQ for RocketMQ. Necessário apenas para instâncias serverless acessadas pela Internet.

Nome de usuário da instância

1XVg0hzgKm******

Página Access Control > aba Intelligent Authentication. Consulte Obter o nome de usuário e a senha de uma instância.

Senha da instância

ijSt8rEc45******

Mesmo local do nome de usuário.

Quando as credenciais são necessárias

A necessidade de fornecer nome de usuário, senha e namespace depende do método de acesso:

Método de acesso

Nome de usuário e senha

Namespace (ID da instância)

VPC (instância padrão)

Não necessário. O broker obtém as credenciais automaticamente das informações da VPC.

Não necessário.

VPC (instância serverless, autenticação gratuita ativada)

Não necessário.

Não necessário.

VPC (instância serverless, autenticação gratuita desativada)

Necessário.

Não necessário.

Internet (instância padrão)

Necessário. Ative primeiro o acesso à Internet na instância.

Não necessário.

Internet (instância serverless)

Necessário.

Necessário. Defina via .setNamespace("InstanceId").

Envie mensagens

Crie um arquivo ProducerExample.java no projeto e execute-o:

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 {
        // Instance endpoint (VPC endpoint for internal access, public endpoint for Internet access)
        String endpoints = "<your-endpoint>";
        // Topic (must be pre-created in the console)
        String topic = "<your-topic>";

        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
                .setEndpoints(endpoints)
                // Uncomment for serverless instances accessed over the Internet:
                // .setNamespace("<your-instance-id>")
                // Uncomment for Internet access or serverless VPC without authentication-free:
                // .setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"))
                .build();

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

        Message message = provider.newMessageBuilder()
                .setTopic(topic)
                .setKeys("messageKey")
                .setTag("messageTag")
                .setBody("messageBody".getBytes())
                .build();
        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println(sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

Substitua os espaços reservados pelos seus valores:

Espaço reservado

Descrição

<your-endpoint>

Endpoint da instância, por exemplo rmq-cn-xxx.cn-hangzhou.rmq.aliyuncs.com:8080

<your-topic>

Nome do tópico, por exemplo normal_test

<your-instance-id>

ID da instância (apenas para acesso serverless via Internet)

<instance-username>

Nome de usuário da instância (acesso via Internet ou VPC serverless sem autenticação gratuita)

<instance-password>

Senha da instância (acesso via Internet ou VPC serverless sem autenticação gratuita)

Nota

.setTopics()

durante a inicialização do producer valida as configurações do tópico antecipadamente. Isso é opcional para mensagens normais (validadas dinamicamente no momento do envio), mas obrigatório para mensagens transacionais para evitar falhas na API de consulta.

Receba mensagens

O ApsaraMQ for RocketMQ fornece dois tipos de consumidor:

Push consumer

Simple consumer

Como funciona

O SDK entrega mensagens a um callback (listener de mensagens).

Sua aplicação busca mensagens e confirma cada uma explicitamente.

Concorrência

Gerenciada pelo SDK.

Gerenciada pela sua aplicação.

Flexibilidade

Menor. O SDK encapsula o fluxo de consumo.

Maior. Operações atômicas permitem criar fluxos de trabalho personalizados.

Mais indicado para

A maioria dos casos de uso em que você precisa processar mensagens recebidas.

Cenários que requerem controle refinado sobre busca, processamento ou confirmação.

Comece com um push consumer, a menos que precise de controle personalizado sobre o fluxo de consumo.

Push consumer

Crie um arquivo PushConsumerExample.java e execute-o:

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);

    private PushConsumerExample() {
    }

    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();
        ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
                .setEndpoints(endpoints)
                // Uncomment for serverless instances accessed over the Internet:
                // .setNamespace("<your-instance-id>")
                // Uncomment for Internet access or serverless VPC without authentication-free:
                // .setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"))
                .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("Consume Message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();

        // Keep the consumer running
        Thread.sleep(Long.MAX_VALUE);

        // To shut down the consumer gracefully, call:
        // pushConsumer.close();
    }
}

Simple consumer

Crie um arquivo SimpleConsumerExample.java e execute-o:

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);

    private SimpleConsumerExample() {
    }

    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();
        ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
                .setEndpoints(endpoints)
                // Uncomment for serverless instances accessed over the Internet:
                // .setNamespace("<your-instance-id>")
                // Uncomment for Internet access or serverless VPC without authentication-free:
                // .setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"))
                .build();

        // Long polling timeout
        Duration awaitDuration = Duration.ofSeconds(10);
        // Subscribe to all messages in the topic (tag = "*")
        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();

        // Maximum messages per pull
        int maxMessageNum = 16;
        // Messages stay invisible to other consumers for this duration after being received
        Duration invisibleDuration = Duration.ofSeconds(10);

        // Poll for messages in a loop. For real-time consumption, use multiple threads.
        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 {
                    // ACK each message to commit the consumption result to the broker
                    consumer.ack(message);
                    System.out.println("Message is acknowledged successfully, messageId= " + messageId);
                } catch (Throwable t) {
                    t.printStackTrace();
                }
            }
        }
        // To shut down the consumer gracefully, call:
        // consumer.close();
    }
}

Verifique a entrega de mensagens

Confira o status de entrega no console do ApsaraMQ for RocketMQ:

  1. Faça login no console do ApsaraMQ for RocketMQ.

  2. Na página Instances, clique no nome da sua instância.

  3. No painel de navegação à esquerda, clique em Message Tracing.

Versões do SDK necessárias para acesso serverless via Internet

O acesso a uma instância serverless do ApsaraMQ for RocketMQ pela Internet exige uma versão mínima do SDK. Substitua InstanceId nos exemplos pelo ID real da sua instância.

SDK for Java 5.x (rocketmq-client-java)

Versão mínima: 5.0.6

Defina o namespace em ClientConfiguration:

ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
    .setEndpoints(endpoints)
    .setNamespace("InstanceId")
    .setCredentialProvider(sessionCredentialsProvider)
    .build();

SDK for Java 5.x (rocketmq-client)

Versão mínima: 5.2.0

Defina o namespace separadamente no producer e no consumer:

// Producer
producer.setNamespaceV2("InstanceId");

// Consumer
consumer.setNamespaceV2("InstanceId");

TCP client SDK for Java 1.x

Versão mínima: 1.9.0.Final

Defina o namespace por meio de propriedades:

properties.setProperty(PropertyKeyConst.Namespace, "InstanceId");

Próximos passos