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 :
Le producteur envoie un demi-message au broker. Le broker stocke ce demi-message, mais ne le distribue pas encore aux consommateurs.
Le broker persiste le demi-message et renvoie un accusé de réception (ACK) au producteur.
Le producteur exécute la transaction locale, par exemple en insérant une commande dans une base de données.
-
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.
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.

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 :
Téléchargé l'édition communautaire du SDK pour Java 4.5.2 ou version ultérieure depuis la page de téléchargement de RocketMQ
Créé une paire AccessKey pour votre compte Alibaba Cloud
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 :
Interroger l'état de la transaction locale correspondant au demi-message, par exemple en recherchant la commande dans votre base de données.
Renvoyer le
LocalTransactionStatecorrect 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();
}
}