Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Sample code for RocketMQ 3.x/4.x SDK for Java

Dernière mise à jour :Aug 09, 2026

Les instances ApsaraMQ for RocketMQ 5.x prennent en charge les clients construits avec les SDK RocketMQ 3.x et 4.x. Cette rubrique fournit des exemples Java pour l'envoi et la réception de messages normaux, ordonnés, planifiés, différés et transactionnels à l'aide de ces SDK hérités.

Important

Les derniers SDK RocketMQ 5.x sont entièrement compatibles avec les brokers ApsaraMQ for RocketMQ 5.x et offrent des fonctionnalités supplémentaires. Utilisez les SDK 5.x pour les nouveaux projets. Pour plus d'informations, consultez les Notes de version. Les SDK client 3.x, 4.x et TCP sont maintenus uniquement pour les charges de travail existantes.

Prérequis

Avant de commencer, assurez-vous que vous disposez des éléments suivants :

  • Une instance ApsaraMQ for RocketMQ 5.x

  • Un topic et un groupe de consommateurs créés dans la console ApsaraMQ for RocketMQ

  • Le SDK RocketMQ 3.x ou 4.x pour Java ajouté à votre projet

  • Le point d'accès de l'instance, au format rmq-cn-XXXX.rmq.aliyuncs.com:8080 (disponible dans la console)

Pour obtenir des instructions complètes sur la configuration, consultez Préparatifs.

Configuration partagée

Tous les exemples de cette rubrique partagent la même configuration d'authentification et de connexion. Les blocs suivants définissent la configuration partagée afin que chaque exemple de type de message puisse se concentrer sur la logique d'envoi ou de réception.

Authentification

ApsaraMQ for RocketMQ prend en charge deux méthodes d'accès. Choisissez celle qui correspond à votre environnement réseau :

  • Point de terminaison public -- Configurez RPCHook avec le nom d'utilisateur et le mot de passe de l'instance. Obtenez ces informations d'identification depuis l'onglet Intelligent Authentication de la page Access Control dans la console ApsaraMQ for RocketMQ. N'utilisez pas la paire AccessKey de votre compte Alibaba Cloud.

  • Point de terminaison VPC -- Aucune information d'identification n'est requise. Lorsque le client s'exécute sur une instance Elastic Compute Service (ECS) au sein d'un Virtual Private Cloud (VPC), le broker obtient automatiquement les informations d'identification à partir de la configuration du VPC.

Pour les instances serverless, les informations d'identification sont requises pour l'accès via le point de terminaison public. Si vous activez la fonctionnalité d'accès sans authentification dans les VPC, l'accès VPC ne nécessite pas d'informations d'identification.

Utilitaire RPCHook

import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.remoting.RPCHook;

/**
 * Returns an RPCHook configured with instance credentials.
 * Required for public endpoint access and serverless instances over the Internet.
 * Not required for VPC access on standard instances.
 */
private static RPCHook getAclRPCHook() {
    return new AclClientRPCHook(new SessionCredentials(
        "<instance-username>",   // From the Intelligent Authentication tab
        "<instance-password>"    // From the Intelligent Authentication tab
    ));
}

Initialisation du producteur

import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.producer.DefaultMQProducer;

// --- Public endpoint: pass RPCHook ---
DefaultMQProducer producer = new DefaultMQProducer(getAclRPCHook());

// --- VPC endpoint: omit RPCHook ---
// DefaultMQProducer producer = new DefaultMQProducer();

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

// Required for message trace. Omit to disable tracing.
producer.setAccessChannel(AccessChannel.CLOUD);

// Access point from the console. Use the domain and port only -- no http:// prefix, no resolved IP.
producer.setNamesrvAddr("<your-access-point>");

producer.start();

Initialisation du consommateur

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;

// --- Public endpoint: pass RPCHook ---
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(getAclRPCHook());

// --- VPC endpoint: omit RPCHook ---
// DefaultMQPushConsumer consumer = new DefaultMQPushConsumer();

consumer.setConsumerGroup("<your-group-id>");
consumer.setAccessChannel(AccessChannel.CLOUD);
consumer.setNamesrvAddr("<your-access-point>");
consumer.subscribe("<your-topic>", "*");

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 provenant de l'onglet Intelligent Authentication MjY1NTI5MD...
<instance-password> Mot de passe de l'instance provenant de l'onglet Intelligent Authentication OTk2QjFBMz...
<your-group-id> ID du groupe de consommateurs créé dans la console GID_test_group
<your-access-point> Point d'accès de l'instance (domaine:port) provenant de la console rmq-cn-hangzhou.rmq.aliyuncs.com:8080
<your-topic> Topic créé dans la console test_topic

Messages normaux

Les messages normaux prennent en charge trois modes d'envoi : synchrone, asynchrone et unidirectionnel. Pour plus de détails conceptuels, consultez Messages normaux.

Envoyer des messages normaux de manière synchrone

L'envoi synchrone bloque jusqu'à ce que le broker accuse réception de chaque message. Utilisez ce mode lorsque la confirmation de livraison est requise.

import java.util.Date;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;

// After producer initialization (see Shared configuration)

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) {
        // Handle send failure: retry or persist the message
        System.out.println(new Date() + " Send mq message failed.");
        e.printStackTrace();
    }
}

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

Envoyer des messages normaux de manière asynchrone

L'envoi asynchrone retourne immédiatement le contrôle et fournit les résultats via un rappel. Utilisez ce mode pour les producteurs sensibles à la latence.

import java.util.Date;
import java.util.concurrent.TimeUnit;
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.common.RemotingHelper;

// After producer initialization (see Shared configuration)

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) {
                // Handle send failure: retry or persist the message
                System.out.println("send message failed.");
                throwable.printStackTrace();
            }
        });
    } catch (Exception e) {
        System.out.println(new Date() + " Send mq message failed.");
        e.printStackTrace();
    }
}

// Wait for async callbacks to complete before shutting down
TimeUnit.SECONDS.sleep(3);
producer.shutdown();

Envoyer des messages normaux en mode unidirectionnel

L'envoi unidirectionnel transmet un message sans attendre de réponse du broker. Utilisez ce mode pour les charges de travail à haut débit où une perte occasionnelle de messages est acceptable, comme la collecte de journaux.

import java.util.Date;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;

// After producer initialization (see Shared configuration)

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

        producer.sendOneway(msg);
    } catch (Exception e) {
        System.out.println(new Date() + " Send mq message failed.");
        e.printStackTrace();
    }
}

producer.shutdown();

Recevoir des messages normaux

Abonnez-vous aux messages normaux avec un écouteur simultané. Chaque message est traité indépendamment et l'ordre n'est pas garanti.

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.common.message.MessageExt;
import java.util.List;

// After consumer initialization (see Shared configuration)

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 selon le principe premier entré, premier sorti (FIFO) au sein de la même file d'attente de messages. Les messages partageant la même clé de partitionnement sont routés vers la même file d'attente. Pour plus de détails conceptuels, consultez Messages ordonnés.

Envoyer des messages ordonnés

Contrairement aux messages normaux, les messages ordonnés nécessitent un MessageQueueSelector pour router les messages ayant la même clé de partitionnement vers la même file d'attente.

import java.util.List;
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.common.RemotingHelper;

// After producer initialization (see Shared configuration)

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 ensure ordered delivery.
        // For ApsaraMQ for RocketMQ 5.x instances, an alternative is:
        // 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) {
                // Route messages with the same orderId to the same queue
                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();

Recevoir des messages ordonnés

Utilisez MessageListenerOrderly au lieu de MessageListenerConcurrently. Le broker livre les messages dans l'ordre au sein de chaque file d'attente.

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.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;

// After consumer initialization (see Shared configuration)
// Subscribe to the ordered topic instead of a normal topic:
// 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 à une heure spécifiée. Les messages différés sont livrés après un délai spécifié. Les deux utilisent la propriété __STARTDELIVERTIME. Pour plus de détails conceptuels, consultez Messages planifiés et différés.

Envoyer des messages planifiés ou différés

import java.util.Date;
import java.text.SimpleDateFormat;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;

// After producer initialization (see Shared configuration)

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

        // Option 1: Delayed delivery -- deliver 3 seconds from now
        long delayTime = System.currentTimeMillis() + 3000;
        msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(delayTime));

        // Option 2: Scheduled delivery -- deliver at a specific time
        // Format: yyyy-MM-dd HH:mm:ss
        // If the time is in the past, 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();

Recevoir des messages planifiés ou différés

Le code du consommateur est identique à celui de la section Recevoir des messages normaux. Aucune configuration spéciale n'est requise côté consommateur.

Messages transactionnels

Les messages transactionnels fournissent une prise en charge des transactions distribuées via un commit en deux phases : le producteur envoie un demi-message, exécute une transaction locale, puis valide ou annule le message en fonction du résultat de la transaction. Pour plus de détails conceptuels, consultez Messages transactionnels.

Important

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

Envoyer des messages transactionnels

Cet exemple envoie un demi-message et exécute une transaction locale. Le broker conserve le demi-message jusqu'à ce que le producteur le valide ou l'annule.

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.common.RemotingHelper;

// --- Public endpoint: pass RPCHook and a dedicated transaction group ID ---
TransactionMQProducer transactionMQProducer = new TransactionMQProducer(
    "<your-transaction-group-id>", getAclRPCHook());

// --- VPC endpoint: omit RPCHook ---
// TransactionMQProducer transactionMQProducer = new TransactionMQProducer(
//     "<your-transaction-group-id>");

transactionMQProducer.setAccessChannel(AccessChannel.CLOUD);
transactionMQProducer.setNamesrvAddr("<your-access-point>");

// Register the transaction status checker (see next section)
transactionMQProducer.setTransactionCheckListener(new LocalTransactionCheckerImpl());
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,
            new LocalTransactionExecuter() {
                @Override
                public LocalTransactionState executeLocalTransactionBranch(
                        Message msg, Object arg) {
                    System.out.println("Start executing the local transaction: " + msg);
                    // Run your local transaction logic here.
                    // Return COMMIT_MESSAGE, ROLLBACK_MESSAGE, or UNKNOW.
                    return LocalTransactionState.UNKNOW;
                }
            },
            null);

        assert sendResult != null;
    } catch (Exception e) {
        e.printStackTrace();
    }
}

Vérifier l'état de la transaction

Lorsque le producteur renvoie UNKNOW, le broker appelle périodiquement checkLocalTransactionState pour résoudre la transaction. Implémentez l'interface TransactionCheckListener pour interroger l'état de votre transaction locale et renvoyer le résultat final.

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("Transaction status check received. MsgId: " + msg.getMsgId());
        // Query your local data store to determine the transaction outcome.
        // Return COMMIT_MESSAGE, ROLLBACK_MESSAGE, or UNKNOW.
        return LocalTransactionState.COMMIT_MESSAGE;
    }
}

Recevoir des messages transactionnels

Le code du consommateur est identique à celui de la section Recevoir des messages normaux. Aucune configuration spéciale n'est requise côté consommateur.