Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Sample code for the RocketMQ 1.x TCP client SDK for Java

Dernière mise à jour :Aug 09, 2026

Les instances ApsaraMQ for RocketMQ 5.x prennent en charge le SDK client TCP RocketMQ 1.x pour Java. Les exemples de code suivants couvrent quatre types de messages : standard, ordonné, planifié/différé et transactionnel.

Important

Les derniers SDK RocketMQ 5.x sont entièrement compatibles avec les brokers 5.x et offrent davantage de fonctionnalités. Utilisez-les pour tous les nouveaux projets. Pour plus d'informations, consultez les Notes de version. Alibaba Cloud ne maintient que les SDK clients TCP RocketMQ 3.x, 4.x. Utilisez-les uniquement pour les charges de travail existantes.

Index du code

Utilisez le tableau suivant pour accéder directement au code dont vous avez besoin.

Type de message Envoi Réception
Standard Synchrone Push
Standard Asynchrone Push par lot
Standard Unidirectionnel Pull
Standard Multithread
Ordonné Envoyer S'abonner
Planifié Envoyer Identique aux messages standards
Différé Envoyer Identique aux messages standards
Transactionnel Envoyer Identique aux messages standards

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Une instance ApsaraMQ for RocketMQ 5.x avec l'endpoint, les topics et les groupes de consommateurs créés dans la console

  • Le SDK client TCP RocketMQ 1.x pour Java ajouté à votre projet. Pour plus d'informations, consultez la section Préparatifs

Configuration commune

Tous les exemples partagent la même configuration de connexion et d'authentification présentée ci-dessous. Chaque bloc de code suivant omet cette configuration répétitive et se concentre sur la logique spécifique au type de message.

import com.aliyun.openservices.ons.api.PropertyKeyConst;
import java.util.Properties;

Properties properties = new Properties();

// Instance credentials
// Obtain the username and password on the Intelligent Authentication tab
// of the Access Control page for your instance in the ApsaraMQ for RocketMQ console.
// Do not use your Alibaba Cloud account AccessKey pair.
properties.put(PropertyKeyConst.AccessKey, "<instance-username>");
properties.put(PropertyKeyConst.SecretKey, "<instance-password>");

// Endpoint
// Enter the domain name and port from the console (e.g., rmq-cn-XXXX.rmq.aliyuncs.com:8080).
// Do not add an http:// or https:// prefix. Do not use a resolved IP address.
properties.put(PropertyKeyConst.NAMESRV_ADDR, "<endpoint>");

// Send timeout in milliseconds
properties.setProperty(PropertyKeyConst.SendMsgTimeoutMillis, "3000");

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

Espace réservé Description Exemple
<instance-username> Nom d'utilisateur de l'instance obtenu depuis la console LTAI5tXxx
<instance-password> Mot de passe de l'instance obtenu depuis la console xXxXxXx
<endpoint> Endpoint de l'instance obtenu depuis la console rmq-cn-XXXX.rmq.aliyuncs.com:8080

Règles d'authentification

Méthode d'accès Identifiants requis ?
Endpoint public Oui. Définissez AccessKey et SecretKey
Virtual Private Cloud (VPC) depuis une instance Elastic Compute Service (ECS) Non. Le broker récupère automatiquement les identifiants à partir des informations VPC
Instance Serverless via Internet Oui
Instance Serverless dans un VPC (authentification sans mot de passe activée) Non
Ne définissez pas l'ID de l'instance lors de l'utilisation d'un SDK client TCP pour se connecter à une instance 5.x. Cela entraînerait des échecs de connexion.

Instances Serverless via Internet

Si vous accédez à une instance serverless via Internet, utilisez la version 1.9.0.Final du SDK ou une version ultérieure, et ajoutez la propriété suivante :

properties.setProperty(PropertyKeyConst.Namespace, "<instance-id>");

Remplacez <instance-id> par l'ID de votre instance ApsaraMQ for RocketMQ.

Envoyer et recevoir des messages standards

Envoyer des messages standards de manière synchrone

L'envoi synchrone bloque l'exécution jusqu'à ce que le broker accuse réception du message. Utilisez ce mode lorsque la confirmation de livraison est cruciale.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.SendResult;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;

import java.util.Date;
import java.util.Properties;

public class ProducerTest {
    public static void main(String[] args) {
        Properties properties = new Properties();
        // ... common configuration (see above) ...

        Producer producer = ONSFactory.createProducer(properties);
        // Call start() once before sending any messages.
        producer.start();

        for (int i = 0; i < 100; i++) {
            Message msg = new Message(
                    "TopicTestMQ",           // Topic created in the console.
                                             // A normal message topic cannot send or receive other message types.
                    "TagA",                  // Tag for consumer-side filtering.
                    "Hello MQ".getBytes());  // Message body (binary). Producer and consumer must
                                             // agree on serialization and deserialization methods.

            // Business key. Should be globally unique when possible.
            // Use it to look up messages in the console if delivery issues occur.
            msg.setKey("ORDERID_" + i);

            try {
                SendResult sendResult = producer.send(msg);
                // No exception means the message was sent successfully.
                if (sendResult != null) {
                    System.out.println(new Date() + " Send mq message success. Topic is:" + msg.getTopic() + " msgId is: " + sendResult.getMessageId());
                }
            } catch (Exception e) {
                // Implement retry or persistence logic here.
                System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
                e.printStackTrace();
            }
        }

        // Shut down the producer before exiting the application.
        // Skip this if the application sends messages continuously.
        producer.shutdown();
    }
}

Envoyer des messages standards de manière asynchrone

L'envoi asynchrone retourne immédiatement le contrôle et délivre le résultat via un callback. Privilégiez ce mode pour un débit plus élevé lorsque le thread appelant ne doit pas être bloqué.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.OnExceptionContext;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.SendCallback;
import com.aliyun.openservices.ons.api.SendResult;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;

import java.util.Properties;
import java.util.concurrent.TimeUnit;

public class ProducerTest {
    public static void main(String[] args) throws InterruptedException {
        Properties properties = new Properties();
        // ... common configuration (see above) ...

        Producer producer = ONSFactory.createProducer(properties);
        producer.start();

        Message msg = new Message(
                "TopicTestMQ",
                "TagA",
                "Hello MQ".getBytes());

        msg.setKey("ORDERID_100");

        // sendAsync returns immediately. The result is delivered through the callback.
        producer.sendAsync(msg, new SendCallback() {
            @Override
            public void onSuccess(final SendResult sendResult) {
                System.out.println("send message success. topic=" + sendResult.getTopic() + ", msgId=" + sendResult.getMessageId());
            }

            @Override
            public void onException(OnExceptionContext context) {
                // Implement retry or persistence logic here.
                System.out.println("send message failed. topic=" + context.getTopic() + ", msgId=" + context.getMessageId());
            }
        });

        // Wait for the async callback before shutting down.
        TimeUnit.SECONDS.sleep(3);

        producer.shutdown();
    }
}

Envoyer des messages standards en mode unidirectionnel

L'envoi unidirectionnel transmet le message sans attendre de réponse du broker. Ce mode offre le débit le plus élevé, mais ne garantit pas la livraison. Utilisez-le uniquement pour les scénarios tolérant une perte occasionnelle de messages, tels que la collecte de journaux.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;

import java.util.Properties;

public class ProducerTest {
    public static void main(String[] args) {
        Properties properties = new Properties();
        // ... common configuration (see above) ...

        Producer producer = ONSFactory.createProducer(properties);
        producer.start();

        for (int i = 0; i < 100; i++) {
            Message msg = new Message(
                    "TopicTestMQ",
                    "TagA",
                    "Hello MQ".getBytes());

            msg.setKey("ORDERID_" + i);

            // sendOneway does not wait for a broker response.
            // If delivery confirmation matters, use synchronous or asynchronous sending instead.
            producer.sendOneway(msg);
        }

        producer.shutdown();
    }
}

Envoyer des messages standards depuis plusieurs threads

Le producteur est thread-safe. Partagez une seule instance de producteur entre les threads plutôt que d'en créer une par thread.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;

import java.util.Date;
import java.util.Properties;

public class SharedProducer {
    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.put(PropertyKeyConst.GROUP_ID, "XXX");
        // ... common configuration (see above) ...

        Producer producer = ONSFactory.createProducer(properties);
        producer.start();

        // Both threads share the same producer instance.
        Thread thread = new Thread(new Runnable() {
            @Override
            public void run() {
                Message msg = new Message(
                        "TopicTestMQ",
                        "TagA",
                        "Hello MQ".getBytes());
                try {
                    SendResult sendResult = producer.send(msg);
                    if (sendResult != null) {
                        System.out.println(new Date() + " Send mq message success. Topic is:" + msg.getTopic() + " msgId is: " + sendResult.getMessageId());
                    }
                } catch (Exception e) {
                    System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
                    e.printStackTrace();
                }
            }
        });
        thread.start();

        Thread anotherThread = new Thread(new Runnable() {
            @Override
            public void run() {
                Message msg = new Message("TopicTestMQ", "TagA", "Hello MQ".getBytes());
                try {
                    SendResult sendResult = producer.send(msg);
                    if (sendResult != null) {
                        System.out.println(new Date() + " Send mq message success. Topic is:" + msg.getTopic() + " msgId is: " + sendResult.getMessageId());
                    }
                } catch (Exception e) {
                    System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
                    e.printStackTrace();
                }
            }
        });
        anotherThread.start();

        // (Optional) Shut down the producer when it is no longer needed.
        // producer.shutdown();
    }
}

S'abonner aux messages standards en mode push

En mode push, le broker livre les messages aux consommateurs dès leur arrivée. Il s'agit du modèle de consommation le plus courant.

import com.aliyun.openservices.ons.api.Action;
import com.aliyun.openservices.ons.api.ConsumeContext;
import com.aliyun.openservices.ons.api.Consumer;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.MessageListener;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;

import java.util.Properties;

public class ConsumerTest {
    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.put(PropertyKeyConst.GROUP_ID, "XXX");
        // ... common configuration (see above) ...

        // Consumption mode (optional):
        // properties.put(PropertyKeyConst.MessageModel, PropertyValueConst.CLUSTERING);    // Default
        // properties.put(PropertyKeyConst.MessageModel, PropertyValueConst.BROADCASTING);

        Consumer consumer = ONSFactory.createConsumer(properties);

        // Subscribe to multiple tags with "||".
        consumer.subscribe("TopicTestMQ", "TagA||TagB", new MessageListener() {
            public Action consume(Message message, ConsumeContext context) {
                System.out.println("Receive: " + message);
                return Action.CommitMessage;
            }
        });

        // Subscribe to another topic. Use "*" to receive all tags.
        // To unsubscribe, remove the subscription code and restart the consumer.
        consumer.subscribe("TopicTestMQ-Other", "*", new MessageListener() {
            public Action consume(Message message, ConsumeContext context) {
                System.out.println("Receive: " + message);
                return Action.CommitMessage;
            }
        });

        consumer.start();
        System.out.println("Consumer Started");
    }
}

S'abonner aux messages standards en mode push par lot

Le mode push par lot délivre plusieurs messages lors de chaque invocation du callback, réduisant ainsi la surcharge liée à chaque message.

import com.aliyun.openservices.ons.api.Action;
import com.aliyun.openservices.ons.api.ConsumeContext;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.batch.BatchConsumer;
import com.aliyun.openservices.ons.api.batch.BatchMessageListener;
import java.util.List;
import java.util.Properties;

import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;

public class SimpleBatchConsumer {

    public static void main(String[] args) {
        Properties consumerProperties = new Properties();
        consumerProperties.put(PropertyKeyConst.GROUP_ID, "XXX");
        // ... common configuration (see above) ...

        // Maximum number of messages per batch. Valid values: 1 to 1024. Default: 32.
        consumerProperties.setProperty(PropertyKeyConst.ConsumeMessageBatchMaxSize, String.valueOf(128));
        // Maximum wait time between batches in seconds. Valid values: 0 to 450. Default: 0.
        consumerProperties.setProperty(PropertyKeyConst.BatchConsumeMaxAwaitDurationInSeconds, String.valueOf(10));

        BatchConsumer batchConsumer = ONSFactory.createBatchConsumer(consumerProperties);
        batchConsumer.subscribe("TopicTestMQ", "TagA", new BatchMessageListener() {

             @Override
            public Action consume(final List<Message> messages, ConsumeContext context) {
                System.out.printf("Batch-size: %d\n", messages.size());
                // Process multiple messages at a time.
                return Action.CommitMessage;
            }
        });

        batchConsumer.start();
        System.out.println("Consumer start success.");

        // Keep the process alive.
        try {
            Thread.sleep(200000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
}

S'abonner aux messages standards en mode pull

En mode pull, le consommateur contrôle le moment où il récupère les messages. Utilisez ce mode lorsque vous avez besoin d'un contrôle précis sur le rythme de consommation.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.PullConsumer;
import com.aliyun.openservices.ons.api.TopicPartition;
import java.util.List;
import java.util.Properties;
import java.util.Set;

public class PullConsumerClient {
    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.setProperty(PropertyKeyConst.GROUP_ID, "GID-xxxxx");
        // ... common configuration (see above) ...

        PullConsumer consumer = ONSFactory.createPullConsumer(properties);
        consumer.start();

        // Get all partitions for the topic.
        Set<TopicPartition> topicPartitions = consumer.topicPartitions("topic-xxx");
        // Assign partitions to pull from.
        consumer.assign(topicPartitions);

        while (true) {
            // Poll with a 3000 ms timeout.
            List<Message> messages = consumer.poll(3000);
            System.out.printf("Received message: %s %n", messages);
        }
    }
}

Envoyer et recevoir des messages ordonnés

Les messages ordonnés garantissent que les messages partageant la même clé de sharding sont consommés dans l'ordre de leur envoi. Utilisez-les pour les scénarios nécessitant un séquencement strict, tels que les mises à jour de statut de commande.

Envoyer des messages ordonnés

Utilisez OrderProducer et transmettez une clé de sharding pour router les messages associés vers la même partition.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;
import com.aliyun.openservices.ons.api.order.OrderProducer;

import java.util.Date;
import java.util.Properties;

public class ProducerClient {

    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.put(PropertyKeyConst.GROUP_ID, "XXX");
        // ... common configuration (see above) ...

        OrderProducer producer = ONSFactory.createOrderProducer(properties);
        producer.start();

        for (int i = 0; i < 1000; i++) {
            String orderId = "biz_" + i % 10;
            Message msg = new Message(
                    "Order_global_topic",              // Topic for ordered messages.
                    "TagA",
                    "send order global msg".getBytes()
            );
            msg.setKey(orderId);

            // Sharding key determines partition routing.
            // Messages with the same sharding key are delivered to the same partition
            // and consumed in order.
            String shardingKey = String.valueOf(orderId);
            try {
                SendResult sendResult = producer.send(msg, shardingKey);
                if (sendResult != null) {
                    System.out.println(new Date() + " Send mq message success. Topic is:" + msg.getTopic() + " msgId is: " + sendResult.getMessageId());
                }
            } catch (Exception e) {
                System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
                e.printStackTrace();
            }
        }

        producer.shutdown();
    }
}

S'abonner aux messages ordonnés

Utilisez OrderConsumer avec un MessageOrderListener. Retournez OrderAction.Suspend en cas d'échec de la consommation afin que le broker retente la livraison du message sans rompre l'ordre.

package com.aliyun.openservices.ons.example.order;

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.order.ConsumeOrderContext;
import com.aliyun.openservices.ons.api.order.MessageOrderListener;
import com.aliyun.openservices.ons.api.order.OrderAction;
import com.aliyun.openservices.ons.api.order.OrderConsumer;

import java.util.Properties;

public class ConsumerClient {

    public static void main(String[] args) {
        Properties properties = new Properties();
        properties.put(PropertyKeyConst.GROUP_ID, "XXX");
        // ... common configuration (see above) ...

        // Retry interval (ms) if consumption fails. Valid values: 10 to 30000.
        properties.put(PropertyKeyConst.SuspendTimeMillis, "100");
        // Maximum retry count for failed messages.
        properties.put(PropertyKeyConst.MaxReconsumeTimes, "20");

        OrderConsumer consumer = ONSFactory.createOrderedConsumer(properties);

        consumer.subscribe(
                "Order_global_topic",
                // Tag filter:
                // "*"           -> subscribe to all tags
                // "TagA||TagB"  -> subscribe to TagA or TagB
                "*",
                new MessageOrderListener() {
                    @Override
                    public OrderAction consume(Message message, ConsumeOrderContext context) {
                        System.out.println(message);
                        // Return OrderAction.Suspend if consumption fails or an exception occurs.
                        return OrderAction.Success;
                    }
                });

        consumer.start();
    }
}

Envoyer et recevoir des messages planifiés

Les messages planifiés spécifient une heure de livraison absolue au lieu d'être envoyés immédiatement. Utilisez-les pour les tâches qui doivent s'exécuter à un moment précis, comme l'envoi d'un rappel 24 heures après l'inscription d'un utilisateur.

Envoyer des messages planifiés

Définissez startDeliverTime sur un horodatage Unix en millisecondes. Le broker livrera le message à cet instant. Si l'horodatage est dans le passé, le broker livrera le message immédiatement.

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;

import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Properties;

public class ProducerDelayTest {
    public static void main(String[] args) {
        Properties properties = new Properties();
        // ... common configuration (see above) ...

        Producer producer = ONSFactory.createProducer(properties);
        producer.start();

        Message msg = new Message(
                "Topic",
                "tag",
                "Hello MQ".getBytes());
        msg.setKey("ORDERID_100");

        try {
            // Deliver at a specific time. Replace with your target timestamp.
            long timeStamp = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").parse("2016-03-07 16:21:00").getTime();
            msg.setStartDeliverTime(timeStamp);

            SendResult sendResult = producer.send(msg);
            System.out.println("Message Id:" + sendResult.getMessageId());
        } catch (Exception e) {
            System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
            e.printStackTrace();
        }

        producer.shutdown();
    }
}

S'abonner aux messages planifiés

Abonnez-vous aux messages planifiés de la même manière qu'aux messages standards. Consultez la section S'abonner aux messages standards en mode push.

Envoyer et recevoir des messages différés

Les messages différés spécifient un délai relatif par rapport à l'heure actuelle, plutôt qu'une heure de livraison absolue. Utilisez-les pour des scénarios tels que la nouvelle tentative d'une opération ayant échoué après une période de temporisation.

Envoyer des messages différés

Définissez startDeliverTime sur System.currentTimeMillis() plus le délai souhaité en millisecondes. Le délai maximal est de 40 jours (3 456 000 000 ms).

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;

import java.util.Date;
import java.util.Properties;

public class ProducerDelayTest {
    public static void main(String[] args) {
        Properties properties = new Properties();
        // ... common configuration (see above) ...

        Producer producer = ONSFactory.createProducer(properties);
        producer.start();

        Message msg = new Message(
                "Topic",
                "tag",
                "Hello MQ".getBytes());
        msg.setKey("ORDERID_100");

        try {
            // Deliver 3 seconds from now.
            long delayTime = System.currentTimeMillis() + 3000;
            msg.setStartDeliverTime(delayTime);

            SendResult sendResult = producer.send(msg);
            if (sendResult != null) {
                System.out.println(new Date() + " Send mq message success. Topic is:" + msg.getTopic() + " msgId is: " + sendResult.getMessageId());
            }
        } catch (Exception e) {
            System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
            e.printStackTrace();
        }

        producer.shutdown();
    }
}

S'abonner aux messages différés

Abonnez-vous aux messages différés de la même manière qu'aux messages standards. Consultez la section S'abonner aux messages standards en mode push.

Envoyer et recevoir des messages transactionnels

Les messages transactionnels lient une transaction locale et la livraison d'un message de sorte qu'ils réussissent ou échouent ensemble. Le flux de travail est le suivant :

  1. Le broker envoie un demi-message.

  2. Votre application exécute la transaction locale.

  3. Selon le résultat de la transaction, le broker valide ou annule le message.

Si le broker ne reçoit pas de validation ou d'annulation dans les délais, il appelle un callback de vérification pour interroger l'état de la transaction.

Les messages transactionnels nécessitent un groupe de consommateurs dédié. Ne partagez pas ce groupe avec d'autres types de messages.

Envoyer des messages transactionnels

Implémentez deux interfaces :

  • LocalTransactionExecuter : Exécute votre transaction locale après que le broker a accepté le demi-message.

  • LocalTransactionChecker : Répond aux vérifications d'état du broker lorsque le résultat de la transaction est inconnu.

package com.aliyun.openservices.tcp.example.producer;

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;
import com.aliyun.openservices.ons.api.exception.ONSClientException;
import com.aliyun.openservices.ons.api.transaction.LocalTransactionChecker;
import com.aliyun.openservices.ons.api.transaction.LocalTransactionExecuter;
import com.aliyun.openservices.ons.api.transaction.TransactionProducer;
import com.aliyun.openservices.ons.api.transaction.TransactionStatus;

import java.util.Date;
import java.util.Properties;

public class SimpleTransactionProducer {

    public static void main(String[] args) {

        Properties properties = new Properties();
        // Transactional messages require a dedicated group ID.
        properties.put(PropertyKeyConst.GROUP_ID, "XXX");
        // ... common configuration (see above) ...

        // Register the transaction status checker before creating the producer.
        LocalTransactionCheckerImpl localTransactionChecker = new LocalTransactionCheckerImpl();
        TransactionProducer transactionProducer = ONSFactory.createTransactionProducer(properties, localTransactionChecker);
        transactionProducer.start();

        Message msg = new Message("XXX", "TagA", "Hello MQ transaction===".getBytes());

        for (int i = 0; i < 3; i++) {
            try {
                SendResult sendResult = transactionProducer.send(msg, new LocalTransactionExecuter() {
                    @Override
                    public TransactionStatus execute(Message msg, Object arg) {
                        // Run the local transaction here.
                        System.out.println("Execute the local transaction and commit the transaction status.");
                        return TransactionStatus.CommitTransaction;
                    }
                }, null);
                assert sendResult != null;
            } catch (ONSClientException e) {
                System.out.println(new Date() + " Send mq message failed! Topic is:" + msg.getTopic());
                e.printStackTrace();
            }
        }

        System.out.println("Send transaction message success.");
    }
}

// Transaction status checker.
// The broker calls check() when it does not receive a commit or rollback.
class LocalTransactionCheckerImpl implements LocalTransactionChecker {

    @Override
    public TransactionStatus check(Message msg) {
        System.out.println("The request to check the transaction status of the message is received. MsgId: " + msg.getMsgID());
        return TransactionStatus.CommitTransaction;
    }
}

S'abonner aux messages transactionnels

Abonnez-vous aux messages transactionnels de la même manière qu'aux messages standards. Consultez la section S'abonner aux messages standards en mode push.

Étapes suivantes

  • Configurez la journalisation côté client pour aider à résoudre les problèmes de livraison des messages. Consultez la section Configurations des journaux.

  • Migrez vers les derniers SDK RocketMQ 5.x pour bénéficier de plus de fonctionnalités et d'une meilleure compatibilité. Consultez les Notes de version.