Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive transactional messages

Dernière mise à jour :Aug 09, 2026

Les messages transactionnels garantissent l'atomicité entre une transaction locale et la distribution des messages en aval. Si la transaction locale est validée, le message est distribué aux consommateurs. En cas d'annulation (rollback), le message est ignoré. Cette approche prévient deux types d'échec : la distribution d'un message pour une transaction non finalisée, et la validation d'une transaction sans distribution du message associé.

Les exemples suivants utilisent le SDK client TCP pour Java (édition communautaire) pour implémenter un producteur et un consommateur de messages transactionnels.

Fonctionnement des messages transactionnels

Un message transactionnel traverse cinq étapes entre le producteur et le broker ApsaraMQ for RocketMQ :

  1. Le producteur envoie un demi-message au broker. Le broker stocke ce demi-message, mais ne le distribue pas encore aux consommateurs.

  2. Le broker persiste le demi-message et renvoie un accusé de réception (ACK) au producteur.

  3. Le producteur exécute la transaction locale, par exemple en insérant une commande dans une base de données.

  4. Le producteur communique le résultat de la transaction au broker. Selon l'issue de la transaction locale, le producteur envoie l'un des trois statuts suivants :

    • LocalTransactionState.COMMIT_MESSAGE : la transaction locale a réussi. Le broker distribue le message aux consommateurs.

    • LocalTransactionState.ROLLBACK_MESSAGE : la transaction locale a échoué. Le broker ignore le message.

    • LocalTransactionState.UNKNOW : le résultat est incertain. Le broker effectue une vérification ultérieure.

  5. Le broker lance une vérification du statut de la transaction s'il ne reçoit ni validation ni annulation définitive. Cela se produit lorsque le producteur renvoie UNKNOW, plante avant de confirmer un statut ou perd sa connectivité réseau. Le broker appelle périodiquement l'écouteur de vérification de transaction du producteur jusqu'à recevoir une validation ou une annulation.

Interaction process between the producer, broker, and transaction check listener

Pour plus d'informations sur le modèle de messages transactionnels, consultez la rubrique Messages transactionnels.

Prérequis

Avant de commencer, assurez-vous d'avoir :

Envoyer des messages transactionnels

L'envoi d'un message transactionnel nécessite deux composants :

  • Un producteur qui envoie le demi-message et exécute la transaction locale.

  • Un écouteur de vérification de transaction qui répond aux demandes de vérification de statut provenant du broker.

Étape 1 : Implémenter l'écouteur de vérification de transaction

Le broker appelle cet écouteur pour vérifier le résultat d'une transaction locale. Après l'envoi d'un message transactionnel, ApsaraMQ for RocketMQ appelle l'opération API LocalTransactionChecker pour répondre à la demande de vérification de statut du broker. Votre implémentation doit interroger l'état réel de la transaction, par exemple en vérifiant l'existence d'une commande dans votre base de données, puis renvoyer le statut approprié.

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("Received transaction status check request. MsgId: " + msg.getMsgId());

        // Query your business data to determine the transaction outcome.
        // Example: Check whether the order associated with this message
        // was committed to the database.
        //
        // boolean orderExists = orderService.existsByOrderId(msg.getKeys());
        // if (orderExists) {
        //     return LocalTransactionState.COMMIT_MESSAGE;
        // }
        // return LocalTransactionState.ROLLBACK_MESSAGE;

        return LocalTransactionState.COMMIT_MESSAGE;
    }
}

Étape 2 : Envoyer un demi-message et exécuter la transaction locale

Le producteur envoie un demi-message au broker et exécute immédiatement la transaction locale dans le rappel executeLocalTransactionBranch. Selon la valeur de retour, le broker valide, annule ou planifie une vérification de statut pour le message.

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

Espace réservé Description Exemple
YOUR TRANSACTION GROUP ID L'ID de groupe créé dans la console ApsaraMQ for RocketMQ GID_tx_order
YOUR TRANSACTION TOPIC La rubrique créée dans la console ApsaraMQ for RocketMQ tx_order_topic
YOUR MESSAGE TAG Un tag pour filtrer les messages TagA
http://xxxx.mq-internet.aliyuncs.com:80 L'endpoint issu de la console ApsaraMQ for RocketMQ, au format http://MQ_INST_XXXX.aliyuncs.com:80 http://MQ_INST_1234.mq-internet.aliyuncs.com:80
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.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.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

public class RocketMQTransactionProducer {

    // Load AccessKey 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")));
    }

    public static void main(String[] args) throws MQClientException {
        // Create a transactional producer with message trace enabled.
        // To disable message trace, use:
        // new TransactionMQProducer("YOUR TRANSACTION GROUP ID", getAclRPCHook());
        TransactionMQProducer producer = new TransactionMQProducer(
            null, "YOUR TRANSACTION GROUP ID", getAclRPCHook(), true, null);

        // Set the endpoint from the ApsaraMQ for RocketMQ console.
        producer.setNamesrvAddr("http://xxxx.mq-internet.aliyuncs.com:80");

        // Required for message trace on Alibaba Cloud.
        producer.setAccessChannel(AccessChannel.CLOUD);

        // Register the transaction check listener (see Step 1).
        producer.setTransactionCheckListener(new LocalTransactionCheckerImpl());

        producer.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));

                // Send the half message and execute the local transaction.
                SendResult sendResult = producer.sendMessageInTransaction(
                    message,
                    new LocalTransactionExecuter() {
                        @Override
                        public LocalTransactionState executeLocalTransactionBranch(
                                Message msg, Object arg) {
                            System.out.println("Executing local transaction: " + msg);
                            // Run your business logic here (e.g., insert order into DB).
                            // Return COMMIT_MESSAGE, ROLLBACK_MESSAGE, or UNKNOW
                            // based on the outcome.
                            return LocalTransactionState.UNKNOW;
                        }
                    },
                    null);
                assert sendResult != null;
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }
}

États de transaction

État Signification Comportement du broker
LocalTransactionState.COMMIT_MESSAGE La transaction locale a réussi Distribue le message aux consommateurs
LocalTransactionState.ROLLBACK_MESSAGE La transaction locale a échoué Ignore le message
LocalTransactionState.UNKNOW Le résultat n'est pas encore déterminé Vérifie périodiquement le statut de la transaction via l'écouteur de vérification

Importance de la vérification du statut de transaction

Si le producteur renvoie UNKNOW, plante ou perd sa connectivité après l'envoi du demi-message, le broker ne peut pas déterminer s'il doit distribuer ou ignorer le message. La vérification du statut de transaction résout ce problème en interrogeant périodiquement le producteur jusqu'à recevoir une validation ou une annulation définitive.

Votre implémentation LocalTransactionCheckerImpl doit :

  1. Interroger l'état de la transaction locale correspondant au demi-message, par exemple en recherchant la commande dans votre base de données.

  2. Renvoyer le LocalTransactionState correct au broker.

S'abonner aux messages transactionnels

Abonnez-vous aux messages transactionnels de la même manière qu'aux messages normaux. Le consommateur ignore la transaction : il ne reçoit le message qu'après sa validation par le broker.

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

Espace réservé Description Exemple
YOUR GROUP ID L'ID de groupe de consommateurs créé dans la console ApsaraMQ for RocketMQ GID_consumer_order
YOUR TOPIC La rubrique à laquelle s'abonner tx_order_topic
http://xxxx.mq-internet.aliyuncs.com:80 L'endpoint issu de la console ApsaraMQ for RocketMQ http://MQ_INST_1234.mq-internet.aliyuncs.com:80
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 {

    // Load AccessKey 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")));
    }

    public static void main(String[] args) throws MQClientException {
        // Create a push consumer with message trace enabled.
        // To disable message trace, use:
        // new DefaultMQPushConsumer("YOUR GROUP ID", getAclRPCHook(),
        //     new AllocateMessageQueueAveragely());
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(
            "YOUR GROUP ID", getAclRPCHook(),
            new AllocateMessageQueueAveragely(), true, null);

        // Set the endpoint from the ApsaraMQ for RocketMQ console.
        consumer.setNamesrvAddr("http://xxxx.mq-internet.aliyuncs.com:80");

        // Required for message trace on Alibaba Cloud.
        consumer.setAccessChannel(AccessChannel.CLOUD);

        // Subscribe to all tags (*) on the topic.
        consumer.subscribe("YOUR TOPIC", "*");

        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(
                    List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
                System.out.printf("Received messages: %s %n", msgs);
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        consumer.start();
    }
}