Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Messages transactionnels

Dernière mise à jour :Aug 09, 2026

Les messages transactionnels sont un type de message spécialisé proposé par ApsaraMQ for RocketMQ qui garantit que la validation d'une transaction locale et la distribution d'un message réussissent ou échouent simultanément. Ce mécanisme de validation en deux phases maintient la synchronisation entre un service central et ses consommateurs en aval, sans la surcharge liée au verrouillage des ressources inhérente aux transactions distribuées XA (eXtended Architecture).

Distributed transaction requirements

Utilisez les messages transactionnels dans les situations suivantes :

  • Un système de commande doit mettre à jour sa base de données et notifier atomiquement les services de logistique, de points de fidélité et de panier.

  • Un service de paiement doit enregistrer un débit et publier un événement destiné aux consommateurs du grand livre comptable en aval.

  • La distribution d'un message sans finalisation de la transaction locale (ou l'inverse) placerait le système dans un état incohérent.

Fonctionnement des messages transactionnels

Pourquoi les messages classiques sont insuffisants

L'association d'une transaction de base de données locale à l'envoi d'un message classique crée une faille permettant à l'une de réussir tandis que l'autre échoue :

  • Le message est envoyé, mais la transaction locale échoue. Les consommateurs en aval agissent sur une modification non validée.

  • La transaction locale est validée, mais l'envoi du message échoue. Les consommateurs en aval ne sont jamais informés de la modification.

  • Un délai d'expiration se produit et ni le producteur ni le broker ne peuvent déterminer s'il faut valider ou annuler la transaction.

Normal message solution

Pourquoi les transactions XA sont trop coûteuses

Le protocole XA permet de coordonner des transactions distribuées entre plusieurs systèmes, mais il verrouille les ressources pendant toute la durée de la transaction. À mesure que le nombre de systèmes participants augmente, la contention des verrous s'accroît et le débit diminue.

Validation en deux phases avec des messages semi-validés

Les messages transactionnels ApsaraMQ for RocketMQ utilisent un protocole de validation en deux phases qui évite ces deux problèmes :

Transactional message solution

  1. Envoyez un message semi-validé. Le producteur envoie un message au broker. Le broker persiste le message et renvoie un accusé de réception (ACK) au producteur. Le message est marqué comme non prêt pour la distribution ; on appelle ce type de message un « message semi-validé » (half message). Les consommateurs en aval ne peuvent pas encore le voir.

  2. Exécutez la transaction locale. Le producteur exécute son opération de base de données locale (par exemple, mise à jour du statut d'une commande de impayé à payé).

  3. Validez ou annulez. Le producteur communique le résultat de la transaction locale au broker :

    • Commit (Valider) : Le broker marque le message semi-validé comme prêt pour la distribution et le distribue aux consommateurs.

    • Rollback (Annuler) : Le broker supprime le message semi-validé. Les consommateurs ne le reçoivent jamais.

  4. Vérifiez l'état de la transaction (récupération). Si le broker ne reçoit aucun résultat de validation ou d'annulation, en raison d'une défaillance réseau ou d'un redémarrage du producteur, il envoie une requête d'état à une instance de producteur dans le cluster. Le producteur vérifie le résultat de la transaction locale et le retransmet au broker.

Transaction status check workflow

Remarque

Pour connaître l'intervalle de requête et le nombre maximal de tentatives, consultez la section Limites des paramètres.

Cycle de vie des messages

Un message transactionnel traverse les états suivants :

Transactional message lifecycle

État Description
Initialisation Le producteur construit le message semi-validé et se prépare à l'envoyer au broker.
Transaction à valider Le broker stocke le message semi-validé dans le système de stockage des transactions. Contrairement à un message classique, le message semi-validé n'est pas persisté par le broker de manière standard. Le message est invisible pour les consommateurs.
Validé pour la consommation La transaction locale aboutit. Le broker stocke le message semi-validé dans le système de stockage, ce qui le rend visible pour les consommateurs.
Annulation du message La transaction locale échoue. Le broker supprime le message semi-validé. Le flux de travail prend fin.
Consommation en cours Un consommateur récupère le message et commence à le traiter. Si le consommateur ne renvoie aucun résultat dans le délai d'expiration configuré, ApsaraMQ for RocketMQ retente la distribution. Pour plus de détails, consultez la section Nouvelles tentatives de consommation.
Validation du résultat de consommation Le consommateur valide le résultat de la consommation. Le message est marqué comme consommé, mais n'est pas immédiatement supprimé.
Suppression du message La période de rétention des messages expire ou l'espace de stockage devient insuffisant. ApsaraMQ for RocketMQ supprime les messages les plus anciens de manière rotative. Consultez la section Stockage et nettoyage des messages.

Par défaut, ApsaraMQ for RocketMQ conserve tous les messages. Un message consommé n'est pas supprimé immédiatement ; les consommateurs peuvent le reconsommer jusqu'à l'expiration de la période de rétention ou jusqu'à la libération de l'espace de stockage.

Envoi de messages transactionnels (Java)

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Un topic dont le paramètre MessageType est défini sur Transaction dans la console ApsaraMQ for RocketMQ.

  • L'endpoint de l'instance (disponible dans l'onglet Endpoints de la page Instance Details)

  • (Le cas échéant) Le nom d'utilisateur et le mot de passe de l'instance (disponibles dans l'onglet Intelligent Authentication de la page Access Control)

Différences par rapport aux messages classiques

L'envoi d'un message transactionnel diffère de l'envoi d'un message classique sur deux points :

  • Vérificateur de transaction requis. Enregistrez un vérificateur de transaction lors de la création du producteur. Ce vérificateur s'exécute automatiquement si le broker interroge l'état de la transaction après un échec.

  • Liaison au topic requise. Liez le topic cible au producteur au moment de la création afin que le vérificateur intégré puisse récupérer l'état de la transaction.

Exemple de code

Créez un producteur avec un vérificateur de transaction, démarrez une transaction, envoyez un message semi-validé, exécutez la transaction locale, puis validez ou annulez la transaction.

Code d'exemple

import java.time.Duration;
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;
import org.apache.rocketmq.client.apis.producer.Transaction;
import org.apache.rocketmq.client.apis.producer.TransactionResolution;
import org.apache.rocketmq.client.java.message.MessageBuilderImpl;
import org.apache.rocketmq.client.apis.message.MessageBuilder;
import org.apache.rocketmq.shaded.com.google.common.base.Strings;

public class ProducerTransactionMessageExample {

    // Simulates checking whether the order exists in the database.
    private static boolean checkOrderById(String orderId) {
        return true;
    }

    // Simulates the local transaction (e.g., inserting an order record).
    private static boolean doLocalTransaction() {
        return true;
    }

    public static void main(String[] args) throws ClientException {
        // Replace with your instance endpoint.
        // Find this on the Endpoints tab of the Instance Details page
        // in the ApsaraMQ for RocketMQ console.
        String endpoints = "<your-instance-endpoint>";

        // The topic must have MessageType set to Transaction.
        String topic = "<your-transaction-topic>";

        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder()
            .setEndpoints(endpoints);

        // Authentication:
        // - Public endpoint: specify username and password (find these on the
        //   Intelligent Authentication tab of the Access Control page).
        // - VPC endpoint on ECS: no credentials needed; the broker resolves
        //   them from the VPC.
        // - Serverless instance: always specify credentials, regardless of
        //   access method.
        builder.setCredentialProvider(
            new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        builder.setRequestTimeout(Duration.ofMillis(5000));
        ClientConfiguration configuration = builder.build();

        MessageBuilder messageBuilder = new MessageBuilderImpl();

        // Build the producer with a transaction checker.
        // The checker runs when the broker queries a half message whose
        // commit/rollback result was not received.
        Producer producer = provider.newProducerBuilder()
            .setTransactionChecker(messageView -> {
                // Look up the order ID attached to the half message.
                // If the order exists in the database, the local transaction
                // committed successfully. Otherwise, roll back.
                final String orderId = messageView.getProperties().get("OrderId");
                if (Strings.isNullOrEmpty(orderId)) {
                    return TransactionResolution.ROLLBACK;
                }
                return checkOrderById(orderId)
                    ? TransactionResolution.COMMIT
                    : TransactionResolution.ROLLBACK;
            })
            .setTopics(topic)
            .setClientConfiguration(configuration)
            .build();

        // Step 1: Begin a transaction.
        final Transaction transaction;
        try {
            transaction = producer.beginTransaction();
        } catch (ClientException e) {
            e.printStackTrace();
            return;
        }

        // Step 2: Build and send the half message.
        Message message = messageBuilder.setTopic(topic)
            .setKeys("messageKey1")
            .setTag("messageTag")
            // Attach a business ID for the transaction checker to query later.
            .addProperty("OrderId", "xxx")
            .setBody("messageBody".getBytes())
            .build();

        final SendReceipt sendReceipt;
        try {
            sendReceipt = producer.send(message, transaction);
        } catch (ClientException e) {
            // Half message send failed. The transaction ends here.
            return;
        }

        // Step 3: Run the local transaction.
        boolean localTransactionOk = doLocalTransaction();

        // Step 4: Commit or roll back based on the local transaction result.
        if (localTransactionOk) {
            try {
                transaction.commit();
            } catch (ClientException e) {
                // If the commit call fails, the broker will invoke the
                // transaction checker to resolve the status.
                e.printStackTrace();
            }
        } else {
            try {
                transaction.rollback();
            } catch (ClientException e) {
                // Log the error. The broker will invoke the transaction
                // checker to resolve the status.
                e.printStackTrace();
            }
        }
    }
}

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

Espace réservé Description Exemple
<your-instance-endpoint> Endpoint de l'instance (disponible dans l'onglet Endpoints de la page Instance Details) xxx-hangzhou.rmq.aliyuncs.com:8080
<your-transaction-topic> Nom du topic dont le paramètre MessageType est défini sur Transaction order-tx-topic
<your-username> Nom d'utilisateur de l'instance (disponible dans l'onglet Intelligent Authentication de la page Access Control) MjoxODgwNzcwODY5MD****
<your-password> Mot de passe de l'instance NEh6cm9FVUl****

Pour obtenir des exemples complets de SDK dans toutes les langues prises en charge, consultez la section SDK Apache RocketMQ 5.x.

Bonnes pratiques

Minimiser les résultats de transaction inconnus

La vérification de l'état de la transaction sert de filet de sécurité en cas d'échec lors de la validation ou de l'annulation. Un volume élevé de vérifications d'état dégrade les performances du système et retarde la distribution des messages. Concevez vos transactions locales de manière à renvoyer un résultat définitif (Commit ou Rollback) aussi rapidement que possible.

Gérer correctement les transactions en cours

Lorsque le broker interroge l'état d'un message semi-validé et que la transaction locale est toujours en cours d'exécution, renvoyez Unknown et non Commit ou Rollback. Le renvoi d'un résultat prématuré peut entraîner une incohérence des données.

Si les vérifications d'état arrivent trop tôt parce que la transaction locale est lente, envisagez les approches suivantes :

  • Augmentez le délai de la première vérification. Configurez un intervalle plus long avant que le broker n'envoie sa première requête d'état. Inconvénient : cela retarde également la récupération des transactions ayant réellement échoué.

  • Détectez explicitement l'état « en cours ». Concevez la logique de la transaction locale de manière à distinguer clairement l'état « toujours en cours » de l'état « échoué », afin que le vérificateur renvoie le statut correct.

Limites

Contrainte Détails
Type de topic Les messages transactionnels nécessitent un topic dont le paramètre MessageType est défini sur Transaction.
Un seul SendReceipt par transaction Chaque transaction ne prend en charge qu'un seul SendReceipt.
Cohérence à terme uniquement Les messages transactionnels garantissent la cohérence entre la transaction locale et la distribution du message. Ils ne garantissent pas la cohérence en temps réel entre les consommateurs en aval. Tant que le message n'est pas distribué, l'état en aval peut être en retard par rapport à la transaction amont. Utilisez les messages transactionnels uniquement lorsque le traitement asynchrone en aval est acceptable.
Responsabilité côté consommateur ApsaraMQ for RocketMQ garantit la distribution des messages validés, mais chaque consommateur en aval doit gérer correctement le traitement. Implémentez une logique de nouvelles tentatives de consommation pour gérer les échecs temporaires. Consultez la section Nouvelles tentatives de consommation.
Délai d'expiration de la transaction Si le broker ne peut pas déterminer le résultat de la transaction après le délai d'expiration configuré et le nombre maximal de tentatives, il annule par défaut le message semi-validé. Consultez la section Limites des paramètres.