ApsaraMQ for RocketMQ prend en charge trois modes de transmission pour les messages normaux. Chaque mode offre un compromis différent entre la fiabilité de la livraison et le débit.
| Mode | Réponse du broker | Fiabilité | Débit | Idéal pour |
|---|---|---|---|---|
| Synchrone | Bloque jusqu'à l'accusé de réception | Aucune perte de message | Élevé | Notifications par e-mail, confirmations d'inscription, messages promotionnels |
| Asynchrone | Livré via un rappel | Aucune perte de message | Élevé | Flux de travail sensibles au temps de réponse, pipelines de transcodage vidéo |
| Unidirectionnel | Aucune | Perte de message possible | Le plus élevé | Collecte de journaux |
Les trois modes partagent la même configuration du producteur. La seule différence réside dans l'appel d'envoi :
Synchrone :
producer.send(msg)-- bloque et renvoieSendResultAsynchrone :
producer.send(msg, callback)-- renvoie immédiatement, le résultat est livré viaSendCallbackUnidirectionnel :
producer.sendOneway(msg)-- envoi sans attente de confirmation, aucune valeur de retour
Prérequis
Édition communautaire du SDK pour Java 4.5.2 ou version ultérieure
Paire AccessKey créée pour votre compte Alibaba Cloud
Configuration courante du producteur
L'initialisation du producteur ci-dessous s'applique aux trois modes. Remplacez les espaces réservés suivants par vos valeurs :
| Espace réservé | Description | Exemple |
|---|---|---|
YOUR GROUP ID |
ID de groupe créé dans la console ApsaraMQ for RocketMQ | GID_example |
YOUR ACCESS POINT |
Endpoint de la console ApsaraMQ for RocketMQ | http://MQ_INST_XXXX.aliyuncs.com:80 |
YOUR TOPIC |
Topic créé dans la console ApsaraMQ for RocketMQ | topic_example |
YOUR MESSAGE TAG |
Tag pour le filtrage des messages | tag_example |
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;
// Retrieve credentials from environment variables
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
// Create a producer with message trace enabled.
// To disable message trace, use: new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook());
DefaultMQProducer producer = new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook(), true, null);
// Set the access channel to CLOUD for message trace on Alibaba Cloud.
// Skip this line if message trace is not needed.
producer.setAccessChannel(AccessChannel.CLOUD);
// Set the endpoint from the ApsaraMQ for RocketMQ console.
producer.setNamesrvAddr("YOUR ACCESS POINT");
producer.start();
Transmission synchrone
Le producteur envoie un message et bloque jusqu'à ce que le broker renvoie une réponse. Cela garantit la confirmation de la livraison avant l'envoi du message suivant.

Compromis : Fiable -- aucun message n'est perdu. Cependant, le débit est inférieur à celui du mode asynchrone car le producteur attend chaque réponse avant d'envoyer le message suivant.
Cas d'utilisation : Notifications par e-mail, confirmations d'inscription, messages promotionnels -- scénarios où la confirmation de la livraison est requise avant de poursuivre.
Exemple de code
import java.util.Date;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class RocketMQProducer {
/**
* Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
* and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
public static void main(String[] args) throws MQClientException {
// Create a producer with message trace enabled.
// To disable message trace: new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook());
DefaultMQProducer producer = new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook(), true, null);
// Set to CLOUD for message trace on Alibaba Cloud. Skip if not needed.
producer.setAccessChannel(AccessChannel.CLOUD);
// Endpoint from the ApsaraMQ for RocketMQ console (format: http://MQ_INST_XXXX.aliyuncs.com:80).
producer.setNamesrvAddr("YOUR ACCESS POINT");
producer.start();
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) {
// Resend or persist the message if delivery fails.
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
// Shut down the producer before exiting the application. (Optional)
producer.shutdown();
}
}
producer.send(msg) bloque et renvoie un SendResult contenant l'ID du message et le statut d'envoi.
Transmission asynchrone
Le producteur envoie un message et passe immédiatement à la suite sans attendre la réponse du broker. Implémentez un SendCallback pour gérer la réponse du broker de manière asynchrone.

Compromis : Même fiabilité que le mode synchrone (aucune perte de message), mais le producteur ne bloque pas entre les envois. La complexité supplémentaire réside dans l'implémentation de SendCallback pour la gestion des succès et des échecs.
Cas d'utilisation : Flux de travail sensibles au temps de réponse avec des processus aval de longue durée. Par exemple, après le téléchargement d'une vidéo, un rappel déclenche le transcodage, et un autre rappel pousse le résultat du transcodage.
Exemple de code
import java.util.Date;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
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.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class RocketMQAsyncProducer {
/**
* Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
* and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
public static void main(String[] args) throws MQClientException {
// Create a producer with message trace enabled.
// To disable message trace: new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook());
DefaultMQProducer producer = new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook(), true, null);
// Set to CLOUD for message trace on Alibaba Cloud. Skip if not needed.
producer.setAccessChannel(AccessChannel.CLOUD);
// Endpoint from the ApsaraMQ for RocketMQ console (format: http://MQ_INST_XXXX.aliyuncs.com:80).
producer.setNamesrvAddr("YOUR ACCESS POINT");
producer.start();
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) {
// Message delivered successfully.
System.out.println("send message success. msgId= " + result.getMsgId());
}
@Override public void onException(Throwable throwable) {
// Resend or persist the message if delivery fails.
System.out.println("send message failed.");
throwable.printStackTrace();
}
});
} catch (Exception e) {
// Resend or persist the message if delivery fails.
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
// Shut down the producer before exiting the application. (Optional)
producer.shutdown();
}
}
producer.send(msg, new SendCallback() {...}) ne bloque pas. Le résultat est livré via onSuccess ou onException.
Transmission unidirectionnelle
Le producteur envoie un message et renvoie immédiatement le contrôle. Aucune réponse du broker n'est renvoyée et aucun rappel n'est déclenché. Un message peut être envoyé en quelques microsecondes.

Compromis : Débit le plus élevé des trois modes, mais des messages peuvent être perdus car le producteur ne peut pas confirmer la livraison.
Cas d'utilisation : Flux de données volumineux et peu critiques où une perte occasionnelle de messages est acceptable, comme la collecte de journaux.
Exemple de code
import java.util.Date;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class RocketMQOnewayProducer {
/**
* Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
* and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
public static void main(String[] args) throws MQClientException {
// Create a producer with message trace enabled.
// To disable message trace: new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook());
DefaultMQProducer producer = new DefaultMQProducer("YOUR GROUP ID", getAclRPCHook(), true, null);
// Set to CLOUD for message trace on Alibaba Cloud. Skip if not needed.
producer.setAccessChannel(AccessChannel.CLOUD);
// Endpoint from the ApsaraMQ for RocketMQ console (format: http://MQ_INST_XXXX.aliyuncs.com:80).
producer.setNamesrvAddr("YOUR ACCESS POINT");
producer.start();
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) {
// Resend or persist the message if delivery fails.
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
// Shut down the producer before exiting the application. (Optional)
producer.shutdown();
}
}
producer.sendOneway(msg) ne renvoie aucun SendResult et n'invoque aucun rappel -- le message est envoyé selon le principe du meilleur effort.
S'abonner aux messages normaux
Quel que soit le mode de transmission utilisé pour envoyer les messages, consommez-les avec DefaultMQPushConsumer. Le broker pousse les messages vers le consommateur, qui les traite via un écouteur MessageListenerConcurrently.
import java.util.List;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
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.client.consumer.rebalance.AllocateMessageQueueAveragely;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;
public class RocketMQPushConsumer {
/**
* Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
* and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
public static void main(String[] args) throws MQClientException {
// Create a consumer with message trace enabled.
// To disable message trace:
// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("YOUR GROUP ID", getAclRPCHook(), new AllocateMessageQueueAveragely());
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("YOUR GROUP ID", getAclRPCHook(), new AllocateMessageQueueAveragely(), true, null);
// Endpoint of the ApsaraMQ for RocketMQ instance.
consumer.setNamesrvAddr("http://xxxx.mq-internet.aliyuncs.com:80");
// Set to CLOUD for message trace on Alibaba Cloud. Skip if not needed.
consumer.setAccessChannel(AccessChannel.CLOUD);
// Topic created in the ApsaraMQ for RocketMQ console.
consumer.subscribe("YOUR TOPIC", "*");
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();
}
}