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:
Criado os recursos no ApsaraMQ for RocketMQ, incluindo um tópico e um grupo de consumidores (Etapa 2 desta série)
Uma IDE Java, como IntelliJ IDEA ou Eclipse (os exemplos usam o IntelliJ IDEA Ultimate)
Adicione a dependência do SDK
Crie um projeto Java na sua IDE.
-
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 |
|
|
|
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. |
|
|
|
Tópico criado para envio e recebimento de mensagens. Crie-o previamente. Consulte Criar um tópico. |
|
|
|
Grupo de consumidores criado. Crie-o previamente. Consulte Criar um grupo de consumidores. |
|
|
|
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 |
|
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 |
|
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 |
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 |
|
|
Endpoint da instância, por exemplo |
|
|
Nome do tópico, por exemplo |
|
|
ID da instância (apenas para acesso serverless via Internet) |
|
|
Nome de usuário da instância (acesso via Internet ou VPC serverless sem autenticação gratuita) |
|
|
Senha da instância (acesso via Internet ou VPC serverless sem autenticação gratuita) |
.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:
Faça login no console do ApsaraMQ for RocketMQ.
Na página Instances, clique no nome da sua instância.
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
Envie e receba mensagens em outras linguagens de programação. Consulte a Visão geral do SDK.
Saiba mais sobre PushConsumer e SimpleConsumer.