Les instances ApsaraMQ for RocketMQ 5.x prennent en charge les clients construits avec les SDK RocketMQ 3.x et 4.x. Cette rubrique fournit des exemples Java pour l'envoi et la réception de messages normaux, ordonnés, planifiés, différés et transactionnels à l'aide de ces SDK hérités.
Les derniers SDK RocketMQ 5.x sont entièrement compatibles avec les brokers ApsaraMQ for RocketMQ 5.x et offrent des fonctionnalités supplémentaires. Utilisez les SDK 5.x pour les nouveaux projets. Pour plus d'informations, consultez les Notes de version. Les SDK client 3.x, 4.x et TCP sont maintenus uniquement pour les charges de travail existantes.
Prérequis
Avant de commencer, assurez-vous que vous disposez des éléments suivants :
Une instance ApsaraMQ for RocketMQ 5.x
Un topic et un groupe de consommateurs créés dans la console ApsaraMQ for RocketMQ
Le SDK RocketMQ 3.x ou 4.x pour Java ajouté à votre projet
Le point d'accès de l'instance, au format
rmq-cn-XXXX.rmq.aliyuncs.com:8080(disponible dans la console)
Pour obtenir des instructions complètes sur la configuration, consultez Préparatifs.
Configuration partagée
Tous les exemples de cette rubrique partagent la même configuration d'authentification et de connexion. Les blocs suivants définissent la configuration partagée afin que chaque exemple de type de message puisse se concentrer sur la logique d'envoi ou de réception.
Authentification
ApsaraMQ for RocketMQ prend en charge deux méthodes d'accès. Choisissez celle qui correspond à votre environnement réseau :
Point de terminaison public -- Configurez
RPCHookavec le nom d'utilisateur et le mot de passe de l'instance. Obtenez ces informations d'identification depuis l'onglet Intelligent Authentication de la page Access Control dans la console ApsaraMQ for RocketMQ. N'utilisez pas la paire AccessKey de votre compte Alibaba Cloud.Point de terminaison VPC -- Aucune information d'identification n'est requise. Lorsque le client s'exécute sur une instance Elastic Compute Service (ECS) au sein d'un Virtual Private Cloud (VPC), le broker obtient automatiquement les informations d'identification à partir de la configuration du VPC.
Pour les instances serverless, les informations d'identification sont requises pour l'accès via le point de terminaison public. Si vous activez la fonctionnalité d'accès sans authentification dans les VPC, l'accès VPC ne nécessite pas d'informations d'identification.
Utilitaire RPCHook
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.remoting.RPCHook;
/**
* Returns an RPCHook configured with instance credentials.
* Required for public endpoint access and serverless instances over the Internet.
* Not required for VPC access on standard instances.
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(
"<instance-username>", // From the Intelligent Authentication tab
"<instance-password>" // From the Intelligent Authentication tab
));
}
Initialisation du producteur
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
// --- Public endpoint: pass RPCHook ---
DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
// --- VPC endpoint: omit RPCHook ---
// DefaultMQProducer producer = new DefaultMQProducer();
// Group ID created in the console
producer.setProducerGroup("<your-group-id>");
// Required for message trace. Omit to disable tracing.
producer.setAccessChannel(AccessChannel.CLOUD);
// Access point from the console. Use the domain and port only -- no http:// prefix, no resolved IP.
producer.setNamesrvAddr("<your-access-point>");
producer.start();
Initialisation du consommateur
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
// --- Public endpoint: pass RPCHook ---
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(getAclRPCHook());
// --- VPC endpoint: omit RPCHook ---
// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer();
consumer.setConsumerGroup("<your-group-id>");
consumer.setAccessChannel(AccessChannel.CLOUD);
consumer.setNamesrvAddr("<your-access-point>");
consumer.subscribe("<your-topic>", "*");
Référence des espaces réservés
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | description | Exemple |
|---|---|---|
<instance-username> |
Nom d'utilisateur de l'instance provenant de l'onglet Intelligent Authentication | MjY1NTI5MD... |
<instance-password> |
Mot de passe de l'instance provenant de l'onglet Intelligent Authentication | OTk2QjFBMz... |
<your-group-id> |
ID du groupe de consommateurs créé dans la console | GID_test_group |
<your-access-point> |
Point d'accès de l'instance (domaine:port) provenant de la console | rmq-cn-hangzhou.rmq.aliyuncs.com:8080 |
<your-topic> |
Topic créé dans la console | test_topic |
Messages normaux
Les messages normaux prennent en charge trois modes d'envoi : synchrone, asynchrone et unidirectionnel. Pour plus de détails conceptuels, consultez Messages normaux.
Envoyer des messages normaux de manière synchrone
L'envoi synchrone bloque jusqu'à ce que le broker accuse réception de chaque message. Utilisez ce mode lorsque la confirmation de livraison est requise.
import java.util.Date;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// After producer initialization (see Shared configuration)
for (int i = 0; i < 128; i++) {
try {
Message msg = new Message(
"<your-topic>",
"<your-message-tag>",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
SendResult sendResult = producer.send(msg);
System.out.printf("%s%n", sendResult);
} catch (Exception e) {
// Handle send failure: retry or persist the message
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
// Shut down the producer before exiting.
// Skip shutdown if your application sends messages continuously.
producer.shutdown();
Envoyer des messages normaux de manière asynchrone
L'envoi asynchrone retourne immédiatement le contrôle et fournit les résultats via un rappel. Utilisez ce mode pour les producteurs sensibles à la latence.
import java.util.Date;
import java.util.concurrent.TimeUnit;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// After producer initialization (see Shared configuration)
for (int i = 0; i < 128; i++) {
try {
Message msg = new Message(
"<your-topic>",
"<your-message-tag>",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
producer.send(msg, new SendCallback() {
@Override
public void onSuccess(SendResult result) {
System.out.println("send message success. msgId= " + result.getMsgId());
}
@Override
public void onException(Throwable throwable) {
// Handle send failure: retry or persist the message
System.out.println("send message failed.");
throwable.printStackTrace();
}
});
} catch (Exception e) {
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
// Wait for async callbacks to complete before shutting down
TimeUnit.SECONDS.sleep(3);
producer.shutdown();
Envoyer des messages normaux en mode unidirectionnel
L'envoi unidirectionnel transmet un message sans attendre de réponse du broker. Utilisez ce mode pour les charges de travail à haut débit où une perte occasionnelle de messages est acceptable, comme la collecte de journaux.
import java.util.Date;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// After producer initialization (see Shared configuration)
for (int i = 0; i < 128; i++) {
try {
Message msg = new Message(
"<your-topic>",
"<your-message-tag>",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
producer.sendOneway(msg);
} catch (Exception e) {
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
producer.shutdown();
Recevoir des messages normaux
Abonnez-vous aux messages normaux avec un écouteur simultané. Chaque message est traité indépendamment et l'ordre n'est pas garanti.
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
// After consumer initialization (see Shared configuration)
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
System.out.printf("Receive New Messages: %s %n", msgs);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
Messages ordonnés
Les messages ordonnés garantissent une livraison selon le principe premier entré, premier sorti (FIFO) au sein de la même file d'attente de messages. Les messages partageant la même clé de partitionnement sont routés vers la même file d'attente. Pour plus de détails conceptuels, consultez Messages ordonnés.
Envoyer des messages ordonnés
Contrairement aux messages normaux, les messages ordonnés nécessitent un MessageQueueSelector pour router les messages ayant la même clé de partitionnement vers la même file d'attente.
import java.util.List;
import org.apache.rocketmq.client.producer.MessageQueueSelector;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// After producer initialization (see Shared configuration)
for (int i = 0; i < 128; i++) {
try {
int orderId = i % 10;
Message msg = new Message(
"<your-order-topic>",
"<your-message-tag>",
"OrderID188",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
// Set the sharding key to ensure ordered delivery.
// For ApsaraMQ for RocketMQ 5.x instances, an alternative is:
// msg.putUserProperty(MessageConst.PROPERTY_SHARDING_KEY, orderId + "");
msg.putUserProperty("__SHARDINGKEY", orderId + "");
SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
// Route messages with the same orderId to the same queue
Integer id = (Integer) arg;
int index = id % mqs.size();
return mqs.get(index);
}
}, orderId);
System.out.printf("%s%n", sendResult);
} catch (Exception e) {
e.printStackTrace();
}
}
producer.shutdown();
Recevoir des messages ordonnés
Utilisez MessageListenerOrderly au lieu de MessageListenerConcurrently. Le broker livre les messages dans l'ordre au sein de chaque file d'attente.
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
// After consumer initialization (see Shared configuration)
// Subscribe to the ordered topic instead of a normal topic:
// consumer.subscribe("<your-order-topic>", "*");
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs,
ConsumeOrderlyContext context) {
context.setAutoCommit(true);
System.out.printf("%s Receive New Messages: %s %n",
Thread.currentThread().getName(), msgs);
// Return SUSPEND_CURRENT_QUEUE_A_MOMENT on failure to retry
return ConsumeOrderlyStatus.SUCCESS;
}
});
consumer.start();
System.out.printf("Consumer Started.%n");
Messages planifiés et différés
Les messages planifiés sont livrés à une heure spécifiée. Les messages différés sont livrés après un délai spécifié. Les deux utilisent la propriété __STARTDELIVERTIME. Pour plus de détails conceptuels, consultez Messages planifiés et différés.
Envoyer des messages planifiés ou différés
import java.util.Date;
import java.text.SimpleDateFormat;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// After producer initialization (see Shared configuration)
for (int i = 0; i < 128; i++) {
try {
Message msg = new Message(
"<your-topic>",
"<your-message-tag>",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
// Option 1: Delayed delivery -- deliver 3 seconds from now
long delayTime = System.currentTimeMillis() + 3000;
msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(delayTime));
// Option 2: Scheduled delivery -- deliver at a specific time
// Format: yyyy-MM-dd HH:mm:ss
// If the time is in the past, the message is delivered immediately.
// long timeStamp = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
// .parse("2021-08-10 18:45:00").getTime();
// msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(timeStamp));
SendResult sendResult = producer.send(msg);
System.out.printf("%s%n", sendResult);
} catch (Exception e) {
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
producer.shutdown();
Recevoir des messages planifiés ou différés
Le code du consommateur est identique à celui de la section Recevoir des messages normaux. Aucune configuration spéciale n'est requise côté consommateur.
Messages transactionnels
Les messages transactionnels fournissent une prise en charge des transactions distribuées via un commit en deux phases : le producteur envoie un demi-message, exécute une transaction locale, puis valide ou annule le message en fonction du résultat de la transaction. Pour plus de détails conceptuels, consultez Messages transactionnels.
Les messages transactionnels nécessitent un groupe de consommateurs dédié. Ne partagez pas le groupe avec d'autres types de messages.
Envoyer des messages transactionnels
Cet exemple envoie un demi-message et exécute une transaction locale. Le broker conserve le demi-message jusqu'à ce que le producteur le valide ou l'annule.
import org.apache.rocketmq.client.producer.LocalTransactionExecuter;
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// --- Public endpoint: pass RPCHook and a dedicated transaction group ID ---
TransactionMQProducer transactionMQProducer = new TransactionMQProducer(
"<your-transaction-group-id>", getAclRPCHook());
// --- VPC endpoint: omit RPCHook ---
// TransactionMQProducer transactionMQProducer = new TransactionMQProducer(
// "<your-transaction-group-id>");
transactionMQProducer.setAccessChannel(AccessChannel.CLOUD);
transactionMQProducer.setNamesrvAddr("<your-access-point>");
// Register the transaction status checker (see next section)
transactionMQProducer.setTransactionCheckListener(new LocalTransactionCheckerImpl());
transactionMQProducer.start();
for (int i = 0; i < 10; i++) {
try {
Message message = new Message(
"<your-transaction-topic>",
"<your-message-tag>",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
SendResult sendResult = transactionMQProducer.sendMessageInTransaction(
message,
new LocalTransactionExecuter() {
@Override
public LocalTransactionState executeLocalTransactionBranch(
Message msg, Object arg) {
System.out.println("Start executing the local transaction: " + msg);
// Run your local transaction logic here.
// Return COMMIT_MESSAGE, ROLLBACK_MESSAGE, or UNKNOW.
return LocalTransactionState.UNKNOW;
}
},
null);
assert sendResult != null;
} catch (Exception e) {
e.printStackTrace();
}
}
Vérifier l'état de la transaction
Lorsque le producteur renvoie UNKNOW, le broker appelle périodiquement checkLocalTransactionState pour résoudre la transaction. Implémentez l'interface TransactionCheckListener pour interroger l'état de votre transaction locale et renvoyer le résultat final.
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionCheckListener;
import org.apache.rocketmq.common.message.MessageExt;
public class LocalTransactionCheckerImpl implements TransactionCheckListener {
@Override
public LocalTransactionState checkLocalTransactionState(MessageExt msg) {
System.out.println("Transaction status check received. MsgId: " + msg.getMsgId());
// Query your local data store to determine the transaction outcome.
// Return COMMIT_MESSAGE, ROLLBACK_MESSAGE, or UNKNOW.
return LocalTransactionState.COMMIT_MESSAGE;
}
}
Recevoir des messages transactionnels
Le code du consommateur est identique à celui de la section Recevoir des messages normaux. Aucune configuration spéciale n'est requise côté consommateur.