Si votre application exige un traitement séquentiel, comme la mise à jour du statut des commandes, l'appariement des transactions ou la synchronisation incrémentielle des données, utilisez les messages ordonnés. ApsaraMQ for RocketMQ distribue et consomme ces messages dans un ordre strict premier entré, premier sorti (FIFO), garantissant ainsi que les systèmes en aval traitent toujours les événements selon leur séquence d'origine.
Cette rubrique explique comment envoyer et consommer des messages ordonnés à l'aide du SDK client Apache RocketMQ TCP pour Java.
Fonctionnement de l'ordonnancement
ApsaraMQ for RocketMQ prend en charge deux portées d'ordonnancement :
| Portée d'ordonnancement | Comportement | Cas d'utilisation |
|---|---|---|
| Ordonnancement global | Tous les messages d'un topic suivent une seule séquence FIFO. | Un ordre total strict est requis et le débit n'est pas une priorité. |
| Ordonnancement par partition | Les messages sont répartis entre les partitions via une clé de sharding. L'ordre FIFO est maintenu au sein de chaque partition, qui fonctionne de manière indépendante. | Vous avez besoin à la fois d'ordonnancement et de parallélisme, par exemple un ordre par ID de commande ou par ID utilisateur. |
Clés de sharding
Une clé de sharding identifie la partition à laquelle appartient un message. Les messages partageant la même clé de sharding sont toujours livrés dans l'ordre. Choisissez des clés de sharding adaptées à votre logique métier ; par exemple, utilisez les ID de commande afin que toutes les mises à jour de statut d'une même commande soient traitées séquentiellement, tandis que les mises à jour de différentes commandes puissent être traitées en parallèle.
Évitez d'utiliser une seule clé de sharding pour tous les messages. Cela concentrerait tout le trafic dans une seule partition et éliminerait tout parallélisme.
La clé de sharding diffère de la clé classique d'un message. Définissez la clé de sharding via la propriété utilisateur __SHARDINGKEY, ou MessageConst.PROPERTY_SHARDING_KEY pour les SDK version 5.x et ultérieures.
Pour plus de détails sur le modèle de message ordonné, consultez la section Messages ordonnés.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
SDK Apache RocketMQ pour Java version 4.5.2 ou ultérieure (téléchargement)
Un environnement de développement configuré (guide de configuration)
Une paire AccessKey pour votre compte Alibaba Cloud (créer une paire AccessKey)
Envoyer des messages ordonnés
Le broker détermine l'ordre des messages en fonction de la séquence d'envoi par un producteur ou un thread unique. Si plusieurs producteurs ou threads envoient des messages simultanément, l'ordre est déterminé par la séquence de réception du broker, qui peut différer de l'ordre métier souhaité.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<YOUR_GROUP_ID> |
ID du groupe de producteurs | GID_order_producer |
<YOUR_ENDPOINT> |
Endpoint obtenu depuis la console ApsaraMQ for RocketMQ | http://MQ_INST_XXXX.aliyuncs.com:80 |
<YOUR_TOPIC> |
Nom du topic pour les messages ordonnés | order_topic |
<YOUR_TAG> |
Tag de message pour le filtrage | TagA |
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.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
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.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class RocketMQOrderProducer {
private static RPCHook getAclRPCHook() {
// Retrieve AccessKey ID and AccessKey secret from environment variables.
// Set ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
// before running this example.
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, use:
// DefaultMQProducer producer = new DefaultMQProducer("<YOUR_GROUP_ID>", getAclRPCHook());
DefaultMQProducer producer = new DefaultMQProducer("<YOUR_GROUP_ID>", getAclRPCHook(), true, null);
// Required for message trace on Alibaba Cloud.
producer.setAccessChannel(AccessChannel.CLOUD);
// Set the endpoint obtained from the ApsaraMQ for RocketMQ console.
producer.setNamesrvAddr("<YOUR_ENDPOINT>");
producer.start();
for (int i = 0; i < 128; i++) {
try {
int orderId = i % 10;
Message msg = new Message("<YOUR_TOPIC>",
"<YOUR_TAG>",
"OrderID188",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
// Set the sharding key. Messages with the same sharding key
// are delivered to the same partition in FIFO order.
// For SDK 5.x and later, use:
// 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();
}
}
Consommer des messages ordonnés
Remplacez <YOUR_GROUP_ID>, <YOUR_ENDPOINT> et <YOUR_TOPIC> par les mêmes valeurs que celles utilisées dans le producteur.
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.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.client.consumer.rebalance.AllocateMessageQueueAveragely;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;
public class RocketMQOrderConsumer {
private static RPCHook getAclRPCHook() {
// Retrieve AccessKey ID and AccessKey secret from environment variables.
// Set ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
// before running this example.
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 push consumer with message trace enabled.
// To disable message trace, use:
// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("<YOUR_GROUP_ID>", getAclRPCHook(), null);
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("<YOUR_GROUP_ID>",
getAclRPCHook(), new AllocateMessageQueueAveragely(), true, null);
// Required for message trace on Alibaba Cloud.
consumer.setAccessChannel(AccessChannel.CLOUD);
// Set the endpoint obtained from the ApsaraMQ for RocketMQ console.
consumer.setNamesrvAddr("<YOUR_ENDPOINT>");
// Subscribe to ordered messages. Use "*" to receive all tags,
// or specify a tag expression to filter messages.
consumer.subscribe("<YOUR_TOPIC>", "*");
// Start consuming from the earliest available offset.
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// Register an orderly message listener to preserve FIFO order.
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 SUCCESS after processing.
// To suspend and retry on failure, return:
// ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT
return ConsumeOrderlyStatus.SUCCESS;
}
});
consumer.start();
System.out.printf("Consumer Started.%n");
}
}
Bonnes pratiques
Envoi
Utilisez un seul thread producteur pour un ordre strict. Si plusieurs producteurs ou threads envoient des messages simultanément, le broker ne peut pas garantir l'ordre métier souhaité. Lorsque l'ordre strict est primordial, envoyez tous les messages connexes depuis un seul thread.
Choisissez des clés de sharding granulaires. Utilisez des identifiants métier tels que les ID de commande ou les ID utilisateur comme clés de sharding, afin que différents groupes de messages puissent être traités en parallèle. Évitez d'utiliser une seule clé de sharding pour tous les messages, car cela concentrerait tout le trafic dans une seule partition et supprimerait le parallélisme.
Consommation
Gérez correctement les échecs de consommation. En cas d'échec, renvoyez
ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENTau lieu deSUCCESSpour suspendre et retenter la consommation, préservant ainsi l'ordre FIFO au sein de la partition.
Étapes suivantes
Messages ordonnés : Comprenez le modèle de message ordonné, les garanties d'ordonnancement et les limites.
Préparer l'environnement : Configurez votre environnement de développement pour le SDK client TCP RocketMQ.