Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive transactional messages

Dernière mise à jour :Aug 09, 2026

ApsaraMQ for RocketMQ propose un traitement des transactions distribuées similaire à l'architecture X/Open XA, garantissant la cohérence transactionnelle au sein des systèmes distribués. Cette rubrique explique comment envoyer et consommer des messages transactionnels à l'aide du SDK client HTTP pour Java.

Fonctionnement des messages transactionnels

Transactional message interaction flow

Un message transactionnel traverse les étapes suivantes :

  1. Le producteur envoie un demi-message au broker. Un demi-message est un message transactionnel qui n'a pas encore été validé (commit) ni annulé (rollback).

  2. Le producteur exécute la transaction locale.

  3. En fonction du résultat de la transaction locale, le producteur valide ou annule le demi-message.

  4. Si le broker ne reçoit ni validation ni annulation dans un délai imparti, il lance une vérification de l'état de la transaction pour interroger le résultat de la transaction locale.

  5. Une fois validé, le message devient disponible pour les consommateurs.

Pour plus d'informations, consultez la section Messages transactionnels.

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Installé le SDK client HTTP pour Java. Pour plus d'informations, consultez la page Préparer l'environnement.

  • Créé une instance, un topic et un groupe de consommateurs ApsaraMQ for RocketMQ dans la console ApsaraMQ for RocketMQ.

  • Stocké une paire AccessKey Alibaba Cloud dans les variables d'environnement ALIBABA_CLOUD_ACCESS_KEY_ID et ALIBABA_CLOUD_ACCESS_KEY_SECRET.

Envoyer des messages transactionnels

L'exemple de code suivant illustre l'envoi de messages transactionnels et la gestion des demi-messages non validés à l'aide du SDK client HTTP pour Java. Il présente quatre scénarios de transaction : validation immédiate, validation différée, validation conditionnelle après nouvelle tentative et annulation.

import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQTransProducer;
import com.aliyun.mq.http.common.AckMessageException;
import com.aliyun.mq.http.model.Message;
import com.aliyun.mq.http.model.TopicMessage;

import java.util.List;

public class TransProducer {

    static void processCommitRollError(Throwable e) {
        if (e instanceof AckMessageException) {
            AckMessageException errors = (AckMessageException) e;
            System.out.println("Commit/Roll transaction error, requestId is:" + errors.getRequestId() + ", fail handles:");
            if (errors.getErrorMessages() != null) {
                for (String errorHandle :errors.getErrorMessages().keySet()) {
                    System.out.println("Handle:" + errorHandle + ", ErrorCode:" + errors.getErrorMessages().get(errorHandle).getErrorCode()
                            + ", ErrorMsg:" + errors.getErrorMessages().get(errorHandle).getErrorMessage());
                }
            }
        }
    }

    public static void main(String[] args) throws Throwable {
        MQClient mqClient = new MQClient(
                // HTTP endpoint. Find this on the Instance Details page > Endpoints tab in the ApsaraMQ for RocketMQ console.
                "${HTTP_ENDPOINT}",
                // AccessKey ID and secret from environment variables.
	              System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
	              System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
        );

        // Topic for transactional messages. Create the topic in the ApsaraMQ for RocketMQ console.
        // A topic supports only one message type. A topic for normal messages cannot send transactional messages.
        final String topic = "${TOPIC}";
        // Instance ID. If the instance has a namespace, specify the ID; otherwise, set to null or "".
        // Check the namespace on the Instance Details page in the ApsaraMQ for RocketMQ console.
        final String instanceId = "${INSTANCE_ID}";
        // Consumer group ID created in the ApsaraMQ for RocketMQ console.
        final String groupId = "${GROUP_ID}";

        final MQTransProducer mqTransProducer = mqClient.getTransProducer(instanceId, topic, groupId);

        for (int i = 0; i < 4; i++) {
            TopicMessage topicMessage = new TopicMessage();
            topicMessage.setMessageBody("trans_msg");
            topicMessage.setMessageTag("a");
            topicMessage.setMessageKey(String.valueOf(System.currentTimeMillis()));
            // Delay (in seconds) before the first transaction status check. Valid values: 10 to 300.
            // After the first check, the broker rechecks every 10 seconds for up to 24 hours.
            topicMessage.setTransCheckImmunityTime(10);
            topicMessage.getProperties().put("a", String.valueOf(i));

            TopicMessage pubResultMsg = null;
            pubResultMsg = mqTransProducer.publishMessage(topicMessage);
            System.out.println("Send---->msgId is: " + pubResultMsg.getMessageId()
                    + ", bodyMD5 is: " + pubResultMsg.getMessageBodyMD5()
                    + ", Handle: " + pubResultMsg.getReceiptHandle()
            );
            if (pubResultMsg != null && pubResultMsg.getReceiptHandle() != null) {
                if (i == 0) {
                    // Commit the first message immediately using its receipt handle.
                    // The broker tracks half messages by receipt handle for commit/rollback.
                    try {
                        mqTransProducer.commit(pubResultMsg.getReceiptHandle());
                        System.out.println(String.format("MessageId:%s, commit", pubResultMsg.getMessageId()));
                    } catch (Throwable e) {
                        // Commit/rollback fails if the receipt handle has expired (TransCheckImmunityTime elapsed).
                        if (e instanceof AckMessageException) {
                            processCommitRollError(e);
                            continue;
                        }
                    }
                }
            }
        }

        // Start a separate thread to poll and resolve uncommitted half messages.
        Thread t = new Thread(new Runnable() {
            public void run() {
                int count = 0;
                while(true) {
                    try {
                        if (count == 3) {
                            break;
                        }
                        List<Message> messages = mqTransProducer.consumeHalfMessage(3, 3);
                        if (messages == null) {
                            System.out.println("No Half message!");
                            continue;
                        }
                        System.out.println(String.format("Half---->MessageId:%s,Properties:%s,Body:%s,Latency:%d",
                                messages.get(0).getMessageId(),
                                messages.get(0).getProperties(),
                                messages.get(0).getMessageBodyString(),
                                System.currentTimeMillis() - messages.get(0).getPublishTime()));

                        for (Message message : messages) {
                            try {
                                if (Integer.valueOf(message.getProperties().get("a")) == 1) {
                                    // Commit the transactional message.
                                    mqTransProducer.commit(message.getReceiptHandle());
                                    count++;
                                    System.out.println(String.format("MessageId:%s, commit", message.getMessageId()));
                                } else if (Integer.valueOf(message.getProperties().get("a")) == 2
                                        && message.getConsumedTimes() > 1) {
                                    // Commit the transactional message.
                                    mqTransProducer.commit(message.getReceiptHandle());
                                    count++;
                                    System.out.println(String.format("MessageId:%s, commit", message.getMessageId()));
                                } else if (Integer.valueOf(message.getProperties().get("a")) == 3) {
                                    // Roll back the transactional message.
                                    mqTransProducer.rollback(message.getReceiptHandle());
                                    count++;
                                    System.out.println(String.format("MessageId:%s, rollback", message.getMessageId()));
                                } else {
                                    // Check the status next time.
                                    System.out.println(String.format("MessageId:%s, unknown", message.getMessageId()));
                                }
                            } catch (Throwable e) {
                                // Commit/rollback fails if the receipt handle has expired
                                // (TransCheckImmunityTime or consumeHalfMessage timeout elapsed).
                                processCommitRollError(e);
                            }
                        }
                    } catch (Throwable e) {
                        System.out.println(e.getMessage());
                    }
                }
            }
        });

        t.start();

        t.join();

        mqClient.close();
    }

}

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

Espace réservé Description Emplacement
${HTTP_ENDPOINT} Endpoint HTTP de votre instance Page Instance Details > onglet Endpoints > section HTTP Endpoint
${TOPIC} Nom du topic Page Topics de la console ApsaraMQ for RocketMQ
${INSTANCE_ID} ID de l'instance. Définissez sur null ou "" si l'instance ne possède pas de namespace Page Instance Details
${GROUP_ID} ID du groupe de consommateurs Page Groups de la console ApsaraMQ for RocketMQ

Paramètres clés

Paramètre Description Contraintes
TransCheckImmunityTime Délai avant la première vérification de l'état de la transaction, en secondes De 10 à 300
consumeHalfMessage(batchSize, waitSeconds) Interroge les demi-messages non validés Dans cet exemple, batchSize est défini sur 3 et waitSeconds est défini sur 3

Résultats de transaction dans cet exemple

L'exemple de code illustre quatre scénarios de transaction basés sur la propriété de message a :

**Message (valeur de a)** Résultat Description
0 Validation immédiate Validé immédiatement après l'envoi
1 Validation différée Validé lors du traitement du demi-message
2 Validation conditionnelle Validé uniquement après au moins une nouvelle tentative (consumedTimes > 1)
3 Annulation Annulé lors du traitement du demi-message
Si un demi-message n'est ni validé ni annulé après la première vérification d'état, le broker effectue de nouvelles vérifications toutes les 10 secondes pendant un maximum de 24 heures.

S'abonner aux messages transactionnels

Les messages transactionnels sont consommés de la même manière que les messages normaux. Le code ci-dessous utilise le long polling pour consommer des messages par lots et les acquitter à l'aide du SDK client HTTP pour Java.

import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQConsumer;
import com.aliyun.mq.http.common.AckMessageException;
import com.aliyun.mq.http.model.Message;

import java.util.ArrayList;
import java.util.List;

public class Consumer {

    public static void main(String[] args) {
        MQClient mqClient = new MQClient(
                // HTTP endpoint. To obtain the HTTP endpoint, log on to the ApsaraMQ for RocketMQ console. In the left-side navigation pane, click Instances. On the Instances page, click the name of your instance. On the Instance Details page, scroll to the Basic Information section and view the HTTP endpoint on the Endpoints tab.
                "${HTTP_ENDPOINT}",
                // The AccessKey ID that is used for authentication.
                "${ACCESS_KEY}",
                // The AccessKey secret that is used for authentication.
                "${SECRET_KEY}"
        );

        // Topic for transactional messages. Create the topic in the ApsaraMQ for RocketMQ console.
        // A topic supports only one message type. A topic for normal messages cannot consume transactional messages.
        final String topic = "${TOPIC}";
        // Consumer group ID created in the ApsaraMQ for RocketMQ console.
        final String groupId = "${GROUP_ID}";
        // Instance ID. If the instance has a namespace, specify the ID; otherwise, set to null or "".
        // Check the namespace on the Instance Details page in the ApsaraMQ for RocketMQ console.
        final String instanceId = "${INSTANCE_ID}";

        final MQConsumer consumer;
        if (instanceId != null && instanceId != "") {
            consumer = mqClient.getConsumer(instanceId, topic, groupId, null);
        } else {
            consumer = mqClient.getConsumer(topic, groupId);
        }

        // Consume messages in a loop. For production use, consume messages with multiple threads for higher throughput.
        do {
            List<Message> messages = null;

            try {
                // Long polling: if no message is available, the request is held at the broker
                // for the specified duration (3 seconds here) before returning.
                messages = consumer.consumeMessage(
                        3,// The maximum number of messages that can be consumed at a time. In this example, the value is set to 3. The largest value you can set is 16.
                        3// The duration of a long polling cycle. Unit: seconds. In this example, the value is set to 3. The largest value you can set is 30.
                );
            } catch (Throwable e) {
                e.printStackTrace();
                try {
                    Thread.sleep(2000);
                } catch (InterruptedException e1) {
                    e1.printStackTrace();
                }
            }
            // No messages available for consumption.
            if (messages == null || messages.isEmpty()) {
                System.out.println(Thread.currentThread().getName() + ": no new message, continue!");
                continue;
            }

            // Process the consumed messages.
            for (Message message : messages) {
                System.out.println("Receive message: " + message);
            }

            // Acknowledge consumed messages. If the broker does not receive an ACK
            // before the delivery retry interval elapses, the message is delivered again.
            // Each delivery attempt generates a new receipt handle.
            {
                List<String> handles = new ArrayList<String>();
                for (Message message : messages) {
                    handles.add(message.getReceiptHandle());
                }

                try {
                    consumer.ackMessage(handles);
                } catch (Throwable e) {
                    // ACK may fail if the receipt handle has expired.
                    if (e instanceof AckMessageException) {
                        AckMessageException errors = (AckMessageException) e;
                        System.out.println("Ack message fail, requestId is:" + errors.getRequestId() + ", fail handles:");
                        if (errors.getErrorMessages() != null) {
                            for (String errorHandle :errors.getErrorMessages().keySet()) {
                                System.out.println("Handle:" + errorHandle + ", ErrorCode:" + errors.getErrorMessages().get(errorHandle).getErrorCode()
                                        + ", ErrorMsg:" + errors.getErrorMessages().get(errorHandle).getErrorMessage());
                            }
                        }
                        continue;
                    }
                    e.printStackTrace();
                }
            }
        } while (true);
    }
}

Paramètres du consommateur

Paramètre Description Valeurs valides
Taille de lot consumeMessage Nombre maximal de messages renvoyés par requête La valeur maximale que vous pouvez définir est 16
Durée de polling consumeMessage Durée pendant laquelle le broker conserve la requête en l'absence de messages disponibles, en secondes La valeur maximale que vous pouvez définir est 30

Rubriques connexes