Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Sample code

Dernière mise à jour :Aug 09, 2026

Envoyez et recevez des messages avec ApsaraMQ for RocketMQ en utilisant le SDK Java Apache RocketMQ. Cette rubrique couvre à la fois le SDK basé sur le protocole gRPC (rocketmq-client-java) et le SDK basé sur le protocole Remoting (rocketmq-client).

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Une instance ApsaraMQ for RocketMQ avec des topics et des ID de groupe créés dans la console ApsaraMQ for RocketMQ

  • L'endpoint de l'instance (par exemple, rmq-cn-XXXX.rmq.aliyuncs.com:8080)

  • La dépendance SDK requise ajoutée à votre projet. Pour plus de détails sur les versions, reportez-vous au Guide des versions

Important

Si vous utilisez une instance Serverless, consultez les exigences de version du SDK pour l'accès via le réseau public avant de poursuivre.

SDK protocole gRPC

Le SDK rocketmq-client-java communique via le protocole gRPC. Les tableaux suivants répertorient les exemples de code pour chaque type de message, hébergés dans le répertoire des clients Apache RocketMQ.

Pour l'intégration avec Spring Boot, consultez rocketmq-v5-client-spring-boot-samples.

Important

Lors de l'envoi de messages transactionnels avec le SDK protocole gRPC, spécifiez un topic au démarrage du producteur. Sans topic défini, les vérifications de transaction sont retardées. Si le message n'est pas envoyé dans un délai de quatre heures, le message transactionnel partiel risque d'être supprimé.

Envoyer des messages

Type de message Exemple de code
Message normal (synchrone) ProducerNormalMessageExample.java
Message normal (asynchrone) AsyncProducerExample.java
Message ordonné ProducerFifoMessageExample.java
Message planifié ou différé ProducerDelayMessageExample.java
Message transactionnel ProducerTransactionMessageExample.java
Message léger LiteProducerExample.java

Consommer des messages

Type de consommateur Exemple de code
Consommateur Push PushConsumerExample.java
Consommateur simple (synchrone) SimpleConsumerExample.java
Consommateur simple (asynchrone) AsyncSimpleConsumerExample.java
Consommateur push léger LitePushConsumerExample.java

Pour plus d'informations sur les consommateurs Push et les consommateurs simples, consultez la section Types de consommateurs.

SDK protocole Remoting

Le SDK rocketmq-client communique via le protocole Remoting. Cette section fournit du code intégré pour chaque type de message.

Pour l'intégration avec Spring Boot, consultez rocketmq-spring-boot-samples.

Configuration commune

Tous les exemples utilisant le protocole Remoting partagent la même configuration d'authentification et de connexion. Consultez cette section en premier lieu, puis reportez-vous au code spécifique au type de message ci-dessous.

Authentification

La méthode d'initialisation du producteur ou du consommateur dépend de la méthode d'accès :

  • Endpoint public -- Transmettez un RPCHook contenant le nom d'utilisateur et le mot de passe de l'instance. Obtenez le nom d'utilisateur et le mot de passe de l'instance depuis l'onglet Intelligent Identity Recognition de la console Resource Access Management. N'utilisez pas l'AccessKey ID et l'AccessKey secret de votre compte Alibaba Cloud.

      private static RPCHook getAclRPCHook() {
          return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
      }
    
      // Pass the RPCHook when creating the producer
      DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
  • Endpoint VPC -- Aucun RPCHook n'est requis. Le serveur s'authentifie sur la base du VPC :

      DefaultMQProducer producer = new DefaultMQProducer();
  • Instance Serverless -- Transmettez un RPCHook pour l'accès via le réseau public. Si l'accès sans mot de passe via le réseau interne est activé, aucun RPCHook n'est nécessaire.

Connexion

// Group ID created in the ApsaraMQ for RocketMQ console
producer.setProducerGroup("<your-group-id>");

// Endpoint from the ApsaraMQ for RocketMQ console
// Use the domain name and port as shown. Do not add http:// or https://.
// Do not use a resolved IP address.
producer.setNamesrvAddr("<your-access-point>");

Traces de messages (facultatif)

Pour activer le traçage des messages cloud, configurez le canal d'accès et activez les traces :

producer.setAccessChannel(AccessChannel.CLOUD);

// Required for SDK v5.3.0 and later, in addition to setAccessChannel
producer.setEnableTrace(true);

Référence des espaces réservés

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

Espace réservé Description Exemple
<instance-username> Nom d'utilisateur de l'instance issu de la console --
<instance-password> Mot de passe de l'instance issu de la console --
<your-group-id> ID de groupe créé dans la console GID_example
<your-access-point> Endpoint issu de la console rmq-cn-XXXX.rmq.aliyuncs.com:8080
<your-topic> Topic créé dans la console topic_example
<your-order-topic> Topic pour les messages ordonnés order_topic_example
<your-transaction-topic> Topic pour les messages transactionnels transaction_topic_example
<your-transaction-group-id> ID de groupe dédié aux messages transactionnels GID_transaction_example
<your-message-tag> Tag de message pour le filtrage TagA

Messages normaux

Les messages normaux conviennent à la plupart des cas d'utilisation où l'ordre et les garanties transactionnelles ne sont pas requis.

Envoyer un message normal (synchrone)

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.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

import java.util.Date;

public class RocketMQProducer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
        producer.setProducerGroup("<your-group-id>");
        producer.setAccessChannel(AccessChannel.CLOUD);
        producer.setEnableTrace(true);
        producer.setNamesrvAddr("<your-access-point>");
        producer.start();

        for (int i = 0; i < 128; i++) {
            try {
                Message msg = new Message("<your-topic>",
                        "<your-message-tag>",
                        "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
                SendResult sendResult = producer.send(msg);
                System.out.printf("%s%n", sendResult);
            } catch (Exception e) {
                System.out.println(new Date() + " Send mq message failed.");
                e.printStackTrace();
            }
        }

        // Shut down the producer before the application exits.
        // Note: Destroying the producer object saves system memory.
        // To send messages continuously, keep the producer running.
        producer.shutdown();
    }
}

Envoyer un message normal (asynchrone)

L'envoi asynchrone retourne immédiatement le contrôle et invoque un callback lorsque le broker répond.

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.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

import java.util.Date;
import java.util.concurrent.TimeUnit;

public class RocketMQAsyncProducer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException, InterruptedException {
        DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
        producer.setProducerGroup("<your-group-id>");
        producer.setAccessChannel(AccessChannel.CLOUD);
        producer.setEnableTrace(true);
        producer.setNamesrvAddr("<your-access-point>");
        producer.start();

        for (int i = 0; i < 128; i++) {
            try {
                Message msg = new Message("<your-topic>",
                        "<your-message-tag>",
                        "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
                producer.send(msg, new SendCallback() {
                    @Override
                    public void onSuccess(SendResult result) {
                        System.out.println("send message success. msgId= " + result.getMsgId());
                    }

                    @Override
                    public void onException(Throwable throwable) {
                        System.out.println("send message failed.");
                        throwable.printStackTrace();
                    }
                });
            } catch (Exception e) {
                System.out.println(new Date() + " Send mq message failed.");
                e.printStackTrace();
            }
        }
        // Wait 3 seconds for the asynchronous callbacks to complete
        TimeUnit.SECONDS.sleep(3);

        producer.shutdown();
    }
}

Envoyer un message normal (unidirectionnel)

Les envois unidirectionnels n'attendent pas de réponse du broker. Utilisez ce mode pour les scénarios à haut débit où une perte occasionnelle de messages est acceptable, comme la collecte de journaux.

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.DefaultMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

import java.util.Date;

public class RocketMQOnewayProducer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
        producer.setProducerGroup("<your-group-id>");
        producer.setAccessChannel(AccessChannel.CLOUD);
        producer.setEnableTrace(true);
        producer.setNamesrvAddr("<your-access-point>");
        producer.start();

        for (int i = 0; i < 128; i++) {
            try {
                Message msg = new Message("<your-topic>",
                        "<your-message-tag>",
                        "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
                // sendOneway does not return a result or throw an exception on failure
                producer.sendOneway(msg);
            } catch (Exception e) {
                System.out.println(new Date() + " Send mq message failed.");
                e.printStackTrace();
            }
        }

        producer.shutdown();
    }
}

Consommer des messages normaux

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.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;

import java.util.List;

public class RocketMQPushConsumer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(getAclRPCHook());
        consumer.setConsumerGroup("<your-group-id>");
        consumer.setAccessChannel(AccessChannel.CLOUD);
        consumer.setEnableTrace(true);
        consumer.setNamesrvAddr("<your-access-point>");

        // Subscribe to a topic. Use "*" to receive all tags, or specify a tag expression.
        consumer.subscribe("<your-topic>", "*");
        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
                                                            ConsumeConcurrentlyContext context) {
                System.out.printf("Receive New Messages: %s %n", msgs);
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });
        consumer.start();
    }
}

Messages ordonnés

Les messages ordonnés garantissent une livraison FIFO (First In, First Out) pour une même clé de sharding. Utilisez-les pour des scénarios tels que le traitement séquentiel d'événements ou la synchronisation de données en temps réel.

Envoyer un message ordonné

Les messages partageant la même clé de sharding sont livrés à la même file d'attente dans l'ordre. Définissez la propriété __SHARDINGKEY afin que le broker route les messages connexes vers la même file d'attente.

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.DefaultMQProducer;
import org.apache.rocketmq.client.producer.MessageQueueSelector;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

import java.util.List;

public class RocketMQOrderProducer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
        producer.setProducerGroup("<your-group-id>");
        producer.setAccessChannel(AccessChannel.CLOUD);
        producer.setEnableTrace(true);
        producer.setNamesrvAddr("<your-access-point>");
        producer.start();

        for (int i = 0; i < 128; i++) {
            try {
                int orderId = i % 10;
                Message msg = new Message("<your-order-topic>",
                        "<your-message-tag>",
                        "OrderID188",
                        "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));

                // Set the sharding key to route related messages to the same queue.
                // In v5.x, you can replace the following line with:
                // msg.putUserProperty(MessageConst.PROPERTY_SHARDING_KEY, orderId + "");
                msg.putUserProperty("__SHARDINGKEY", orderId + "");

                SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
                    @Override
                    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
                        // Select a queue based on the orderId
                        Integer id = (Integer) arg;
                        int index = id % mqs.size();
                        return mqs.get(index);
                    }
                }, orderId);

                System.out.printf("%s%n", sendResult);
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
        producer.shutdown();
    }
}

Consommer des messages ordonnés

Utilisez MessageListenerOrderly au lieu de MessageListenerConcurrently pour préserver l'ordre de consommation au sein de chaque file d'attente.

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.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;

import java.util.List;

public class RocketMQOrderConsumer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(getAclRPCHook());
        consumer.setConsumerGroup("<your-group-id>");
        consumer.setAccessChannel(AccessChannel.CLOUD);
        consumer.setEnableTrace(true);
        consumer.setNamesrvAddr("<your-access-point>");
        consumer.subscribe("<your-order-topic>", "*");

        consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
        consumer.registerMessageListener(new MessageListenerOrderly() {

            @Override
            public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
                context.setAutoCommit(true);
                System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msgs);
                // Return SUSPEND_CURRENT_QUEUE_A_MOMENT on failure to retry
                return ConsumeOrderlyStatus.SUCCESS;
            }
        });

        consumer.start();
        System.out.printf("Consumer Started.%n");
    }
}

Messages planifiés et différés

Les messages planifiés sont livrés à un moment précis. Les messages différés sont livrés après un délai configurable. Les deux utilisent la propriété utilisateur __STARTDELIVERTIME.

Envoyer un message planifié ou différé

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.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

import java.util.Date;

public class RocketMQDelayProducer {
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());
        producer.setProducerGroup("<your-group-id>");
        producer.setAccessChannel(AccessChannel.CLOUD);
        producer.setEnableTrace(true);
        producer.setNamesrvAddr("<your-access-point>");
        producer.start();

        for (int i = 0; i < 128; i++) {
            try {
                Message msg = new Message("<your-topic>",
                        "<your-message-tag>",
                        "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));

                // Delayed message: deliver after 3 seconds
                long delayTime = System.currentTimeMillis() + 3000;
                msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(delayTime));

                // Scheduled message: deliver at a specific time (format: yyyy-MM-dd HH:mm:ss)
                // If the specified time is earlier than the current time, the message is delivered immediately.
                // long timeStamp = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").parse("2021-08-10 18:45:00").getTime();
                // msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(timeStamp));

                SendResult sendResult = producer.send(msg);
                System.out.printf("%s%n", sendResult);
            } catch (Exception e) {
                System.out.println(new Date() + " Send mq message failed.");
                e.printStackTrace();
            }
        }

        producer.shutdown();
    }
}

Consommer des messages planifiés et différés

La consommation des messages planifiés et différés s'effectue de la même manière que pour les messages normaux. Aucune configuration supplémentaire n'est requise.

Messages transactionnels

Les messages transactionnels utilisent un protocole de validation en deux phases. Le producteur envoie un message transactionnel partiel, exécute une transaction locale, puis valide ou annule l'opération. Le broker vérifie périodiquement les transactions non validées en appelant checkLocalTransaction.

Important

L'ID de groupe utilisé pour les messages transactionnels ne peut pas être partagé avec d'autres types de messages. Créez un ID de groupe dédié aux producteurs transactionnels.

Envoyer un message transactionnel

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.LocalTransactionState;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

public class RocketMQTransactionProducer {

    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials("<instance-username>", "<instance-password>"));
    }

    public static void main(String[] args) throws MQClientException {
        // Use a dedicated group ID for transactional messages
        TransactionMQProducer transactionMQProducer = new TransactionMQProducer("<your-transaction-group-id>", getAclRPCHook());
        transactionMQProducer.setAccessChannel(AccessChannel.CLOUD);
        transactionMQProducer.setEnableTrace(true);
        transactionMQProducer.setNamesrvAddr("<your-access-point>");

        transactionMQProducer.setTransactionListener(new TransactionListener() {

            @Override
            public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
                // Run your local transaction logic here
                System.out.println("Start to execute the local transaction: " + msg);
                return LocalTransactionState.UNKNOW;
            }

            @Override
            public LocalTransactionState checkLocalTransaction(MessageExt msg) {
                // Return the local transaction status when the broker checks back
                System.out.println("Received a transaction check request, MsgId: " + msg.getMsgId());
                return LocalTransactionState.COMMIT_MESSAGE;
            }
        });
        transactionMQProducer.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));
                SendResult sendResult = transactionMQProducer.sendMessageInTransaction(message, null);
                assert sendResult != null;
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }
}

Consommer des messages transactionnels

La consommation des messages transactionnels s'effectue de la même manière que pour les messages normaux. Aucune configuration supplémentaire n'est requise.

Accès au réseau public pour les instances Serverless

Pour accéder à une instance Serverless via le réseau public, votre SDK doit respecter une exigence de version minimale. Ajoutez la configuration de namespace indiquée ci-dessous.

Remarque

Remplacez InstanceId par votre ID d'instance réel.

SDK protocole Remoting (rocketmq-client >= 5.2.0)

Ajoutez le namespace au producteur ou au consommateur :

// Producer
producer.setNamespaceV2("InstanceId");

// Consumer
consumer.setNamespaceV2("InstanceId");

SDK protocole gRPC (rocketmq-client-java >= 5.0.6)

Définissez le namespace dans la ClientConfiguration :

ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
    .setEndpoints(endpoints)
    .setNamespace("InstanceId")
    .setCredentialProvider(sessionCredentialsProvider)
    .build();