Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Ordered messages

Dernière mise à jour :Aug 09, 2026

Les messages ordonnés permettent aux consommateurs de traiter les messages exactement dans l'ordre d'envoi. ApsaraMQ for RocketMQ utilise des groupes de messages pour définir la portée de l'ordonnancement : les messages appartenant au même groupe sont toujours livrés de manière séquentielle, tandis que les messages provenant de groupes différents sont traités indépendamment.

Cas d'utilisation

Utilisez les messages ordonnés lorsque les systèmes en aval doivent traiter les événements exactement dans la séquence où ils se sont produits en amont :

  • Appariement des transactions -- Dans le trading de titres financiers, lorsque plusieurs offres partagent le même prix, la règle « premier arrivé, premier servi » s'applique. Le système de traitement des commandes doit gérer les offres exactement dans l'ordre de placement.

    Trade matchmaking

  • Synchronisation incrémentielle des données -- Lors de la réplication des modifications de base de données (insertions, mises à jour, suppressions) via une file d'attente de messages vers un système de recherche, les opérations doivent être rejouées dans leur ordre d'origine. Une relecture désordonnée génère un état incohérent.

    Normal message -- out-of-order risk

    Ordered messages -- consistent replay

Fonctionnement de l'ordonnancement

L'ordonnancement des messages de bout en bout comporte deux volets : l'ordre de production (côté envoi) et l'ordre de consommation (côté réception). Ces deux conditions doivent être remplies pour garantir un traitement strict selon le principe premier entré, premier sorti (FIFO).

Ordre de production

Pour garantir l'ordre de production, respectez les trois conditions suivantes :

  1. Même groupe de messages -- Définissez le même groupe de messages pour tous les messages qui doivent être ordonnés les uns par rapport aux autres. Les messages appartenant à des groupes différents n'ont aucune relation d'ordre entre eux.

  2. Producteur unique -- Envoyez tous les messages connexes depuis un seul producteur. Même avec le même groupe de messages, les messages provenant de producteurs différents dans des systèmes distincts n'ont pas d'ordre déterministe.

  3. Envoi sériel -- Envoyez les messages de manière séquentielle depuis un seul thread. Bien que le client producteur prenne en charge l'accès multithread, les messages envoyés en parallèle depuis différents threads n'ont pas d'ordre déterministe.

Lorsque ces conditions sont réunies, les messages partageant le même groupe de messages sont stockés dans la même file d'attente, dans l'ordre d'envoi.

Storage logic for ordered messages

Comportement de stockage :

  • Les messages du même groupe sont stockés de manière séquentielle dans la même file d'attente.

  • Les messages de groupes différents peuvent coexister dans la même file d'attente, mais leur ordre relatif n'est pas garanti.

Dans le schéma ci-dessus, le Groupe de messages 1 (G1-M1, G1-M2, G1-M3) et le Groupe de messages 4 (G4-M1, G4-M2) partagent la File d'attente 1. ApsaraMQ for RocketMQ garantit l'ordre au sein de chaque groupe, mais pas entre les groupes.

Ordre de consommation

Pour garantir l'ordre de consommation, deux mécanismes agissent de concert :

  1. Ordre de livraison -- Le SDK et le protocole serveur livrent les messages dans l'ordre de stockage. Respectez scrupuleusement le schéma réception-traitement-accusé de réception. Un traitement asynchrone peut rompre la garantie d'ordre.

    Important

    Avec PushConsumer, les messages sont livrés un par un dans l'ordre de stockage. Avec SimpleConsumer, plusieurs messages peuvent arriver lors d'une seule opération d'extraction (pull) ; votre application doit les traiter de manière séquentielle et accuser réception de chacun d'eux avant d'appeler à nouveau receive pour le même groupe. Pour plus de détails, consultez Types de consommateurs.

  2. Nombre limité de tentatives -- Lorsqu'un message ordonné échoue après le nombre maximal de tentatives, il est ignoré afin de débloquer les messages suivants. Choisissez un nombre de tentatives qui équilibre la fiabilité et le risque de bloquer l'intégralité du groupe. Les tentatives au sein d'un groupe de messages :

    • Ne perturbent pas la garantie d'ordre.

    • N'affectent pas les messages des autres groupes de messages.

    • Bloquent les messages suivants du même groupe jusqu'à ce que le message actuel soit résolu.

    Important

    Pendant qu'un message ordonné ayant échoué fait l'objet de nouvelles tentatives, les messages suivants du même groupe sont bloqués. Ils ne sont livrés qu'après la réussite du message actuel ou l'épuisement de ses tentatives.

Combinaisons d'ordres de production et de consommation

Le respect strict du principe FIFO exige à la fois l'ordre de production et l'ordre de consommation. Toutefois, tous les consommateurs n'ont pas besoin d'une livraison ordonnée. Adaptez la configuration en fonction du débit et des exigences d'ordonnancement :

Ordre de production Ordre de consommation Résultat
Groupe de messages défini ; envoi sériel Ordonné FIFO strict au sein de chaque groupe de messages
Groupe de messages défini ; envoi sériel Concurrent Ordre chronologique approximatif ; débit plus élevé
Aucun groupe de messages ; envoi non ordonné Ordonné Ordonnancement strict au niveau de la file d'attente (correspond à l'ordre de stockage, pas à l'ordre d'envoi)
Aucun groupe de messages ; envoi non ordonné Concurrent Ordre chronologique approximatif

Cycle de vie des messages

Un message ordonné traverse cinq états :

Message lifecycle

  1. Initialisé -- Le producteur construit le message et se prépare à l'envoyer.

  2. Prêt -- Le message arrive sur le broker et devient visible pour les consommateurs.

  3. En cours de traitement -- Un consommateur récupère le message et le traite. Si le broker ne reçoit aucun accusé de réception dans le délai imparti, il retente la livraison. Pour plus de détails, consultez Tentatives de consommation.

  4. Accusé de réception -- Le consommateur valide le résultat. Par défaut, ApsaraMQ for RocketMQ conserve tous les messages. Le broker marque le message comme consommé, mais ne le supprime pas immédiatement.

  5. Supprimé -- À l'expiration de la période de rétention ou lorsque l'espace de stockage vient à manquer, le broker supprime les messages les plus anciens de manière progressive. Avant la suppression, les messages peuvent être reconsummés. Pour plus de détails, consultez Stockage et nettoyage des messages.

Important
  • Un message faisant l'objet d'une nouvelle tentative est traité comme un nouveau message. Le cycle de vie du message d'origine prend fin.

  • Pendant qu'un message ordonné fait l'objet de nouvelles tentatives, les messages suivants du même groupe sont bloqués jusqu'à ce que le message actuel aboutisse.

Limites

  • Les messages ordonnés ne peuvent être envoyés qu'aux rubriques dont le paramètre MessageType est défini sur FIFO. Le type de message doit correspondre au type de rubrique.

Important

Si un groupe de consommateurs est configuré pour une livraison ordonnée, tous les messages consommés par ce groupe sont facturés en tant que messages ordonnés, quel que soit leur type réel. Si un ordonnancement strict n'est pas requis, configurez le groupe pour une livraison concurrente afin de réduire les coûts.

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Une rubrique dont le paramètre MessageType est défini sur FIFO, créée dans la console ApsaraMQ for RocketMQ

  • Un groupe de consommateurs configuré pour une livraison ordonnée (requis uniquement si vous avez besoin d'une consommation ordonnée)

  • Le SDK gRPC RocketMQ 5.x installé dans votre projet

Important

Les messages ordonnés ne peuvent être envoyés qu'aux rubriques dont le paramètre MessageType est défini sur FIFO. Un groupe de consommateurs qui n'est pas configuré en mode de livraison ordonnée livre les messages de manière concurrente. Vérifiez à la fois le type de rubrique et le mode de livraison du groupe de consommateurs avant d'écrire tout code.

Optimiser la concurrence de consommation

Avec le SDK gRPC RocketMQ 5.x, PushConsumer peut distribuer les messages provenant de la même MessageQueue à différents threads en fonction de leur groupe de messages. Plus vos valeurs de groupe de messages sont distinctes, plus le gain de débit est important.

Versions de SDK prises en charge :

SDK Version minimale
Java 5.0.8
C++ 5.0.3
Autres SDK Non pris en charge

Envoyer et consommer des messages ordonnés

Les messages ordonnés nécessitent un groupe de messages à chaque appel d'envoi. Concevez des groupes de messages avec la granularité la plus fine permise par votre activité métier -- par exemple, utilisez un ID de commande ou un ID utilisateur. Cela permet un ordonnancement par entité tout en maximisant le parallélisme entre les différentes entités.

Tous les exemples ci-dessous utilisent Java. Pour des exemples complets de SDK dans d'autres langages, consultez le SDK gRPC RocketMQ 5.x.

Exemple de code

Envoyer des messages ordonnés

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. Get it from the Endpoints tab on the instance details page.
        // Use the VPC endpoint for access from ECS over the internal network.
        // Use the public endpoint for access from a local machine or on-premises data center.
        String endpoints = "<your-endpoint>";
        // Topic name. Create the topic in the console first.
        String topic = "<your-fifo-topic>";
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // For public network access, set the instance username and password.
        // Get them from the Intelligent Authentication tab on the Access Control page.
        // For internal network access from ECS, skip this -- the server authenticates via VPC.
        // For Serverless instances, set credentials for public access.
        // If authentication-free internal access is enabled, skip this for internal access.
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"));
        ClientConfiguration configuration = builder.build();
        Producer producer = provider.newProducerBuilder()
                .setTopics(topic)
                .setClientConfiguration(configuration)
                .build();
        // Build an ordered message with a message group.
        Message message = provider.newMessageBuilder()
                .setTopic(topic)
                .setKeys("messageKey")         // Message key for lookup
                .setTag("messageTag")          // Tag for consumer-side filtering
                .setMessageGroup("fifoGroup001") // Message group -- keep values discrete to avoid hot spots
                .setBody("messageBody".getBytes())
                .build();
        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println(sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

Remplacez les espaces réservés suivants par vos valeurs réelles :

Espace réservé Description Exemple
<your-endpoint> Endpoint de l'instance obtenu depuis la console rmq-cn-xxx.cn-hangzhou.rmq.aliyuncs.com:8080
<your-fifo-topic> Nom de la rubrique FIFO order-events
<instance-username> Nom d'utilisateur de l'instance (accès public uniquement) LTAI5tXxx
<instance-password> Mot de passe de l'instance (accès public uniquement) xXxXxXx

Consommer avec PushConsumer

PushConsumer livre les messages ordonnés un par un dans l'ordre de stockage. Configurez le groupe de consommateurs en mode de livraison ordonnée dans la console avant de démarrer le consommateur. Sinon, les messages seront livrés de manière concurrente.

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 {
        // Instance endpoint. Get it from the Endpoints tab on the instance details page.
        String endpoints = "<your-endpoint>";
        // Topic and consumer group. Create both in the console first.
        String topic = "<your-fifo-topic>";
        String consumerGroup = "<your-consumer-group>";
        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // For public network access, set credentials.
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"));
        ClientConfiguration clientConfiguration = builder.build();
        // Subscribe to all tags.
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);
        PushConsumer pushConsumer = provider.newPushConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .setMessageListener(messageView -> {
                    // Process the message and return the result.
                    System.out.println("Consume Message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();
        Thread.sleep(Long.MAX_VALUE);
        // Close the consumer when no longer needed.
        // pushConsumer.close();
    }
}

Consommer avec SimpleConsumer

SimpleConsumer extrait les messages par lots. Traitez les messages de chaque lot de manière séquentielle et accusez réception de chacun d'eux individuellement.

Important

Pour les messages appartenant au même groupe de messages, si un message précédent n'a pas encore fait l'objet d'un accusé de réception, un nouvel appel à receive ne renvoie pas les messages suivants de ce groupe. Ce mécanisme de verrouillage empêche le traitement hors ordre. Les messages appartenant à d'autres groupes de messages ne sont pas affectés et peuvent toujours être reçus de manière concurrente.

import org.apache.rocketmq.client.apis.ClientConfiguration;
import org.apache.rocketmq.client.apis.ClientConfigurationBuilder;
import org.apache.rocketmq.client.apis.ClientException;
import org.apache.rocketmq.client.apis.ClientServiceProvider;
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.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 {
        // Instance endpoint. Get it from the Endpoints tab on the instance details page.
        String endpoints = "<your-endpoint>";
        // Topic and consumer group. Create both in the console first.
        String topic = "<your-fifo-topic>";
        String consumerGroup = "<your-consumer-group>";
        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // For public network access, set credentials.
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"));
        ClientConfiguration clientConfiguration = builder.build();

        Duration awaitDuration = Duration.ofSeconds(10);
        // Subscribe to all tags.
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);
        SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setAwaitDuration(awaitDuration)                    // Long polling timeout
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .build();
        int maxMessageNum = 16;
        Duration invisibleDuration = Duration.ofSeconds(10);       // Invisibility window for processing
        // Poll for messages in a loop. Use multiple threads for higher throughput.
        while (true) {
            final List<MessageView> messageViewList = consumer.receive(maxMessageNum, invisibleDuration);
            messageViewList.forEach(messageView -> {
                System.out.println(messageView);
                // Acknowledge each message after processing.
                try {
                    consumer.ack(messageView);
                } catch (ClientException e) {
                    e.printStackTrace();
                }
            });
        }
        // Close the consumer when no longer needed.
        // consumer.close();
    }
}

Dépanner les tentatives de consommation

Les tentatives de consommation ordonnée pour PushConsumer ont lieu côté client -- le serveur n'enregistre pas les détails des tentatives. Si une trace de message affiche un résultat de livraison failed, vérifiez les journaux du client consommateur.

Pour connaître le chemin d'accès aux journaux du client, consultez Configuration des journaux.

Recherchez les mots-clés suivants dans les journaux du client :

Message listener raised an exception while consuming messages
Failed to consume fifo message finally, run out of attempt times

Bonnes pratiques

Traiter les messages de manière sérielle, pas par lots

Consommez un message à la fois. La consommation par lots peut rompre l'ordonnancement.

Exemple : Les messages sont envoyés dans l'ordre 1 -> 2 -> 3 -> 4. Lors de la consommation par lots, les messages 2 et 3 sont traités ensemble et échouent. Lors de la nouvelle tentative, les messages 2 et 3 sont tous deux redistribués -- mais le message 2 pourrait être traité à nouveau après que le message 3 a déjà réussi lors d'une autre tentative, ce qui entraîne une consommation hors ordre.

Répartir les groupes de messages pour éviter les points chauds

ApsaraMQ for RocketMQ utilise la valeur du groupe de messages pour déterminer quelle file d'attente côté serveur stockera chaque message. Tous les messages d'un même groupe sont routés vers la même file d'attente. Concentrer trop de messages dans quelques groupes surcharge ces files d'attente, créant des points chauds de stockage et limitant la scalabilité.

Utilisez des clés à grain fin comme groupes de messages -- par exemple, des ID de commande ou des ID utilisateur. Plus vos valeurs de groupe de messages sont distinctes, plus les messages sont répartis uniformément entre les files d'attente. Cela permet de maintenir l'ordre des messages pour une même entité tout en répartissant la charge entre les files d'attente pour différentes entités.