Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive normal messages

Dernière mise à jour :Aug 09, 2026

Cette rubrique propose des exemples Java pour envoyer et recevoir des messages normaux via le SDK client HTTP d'ApsaraMQ for RocketMQ. Les messages normaux constituent le type par défaut : ils ne comportent aucune sémantique de livraison particulière, contrairement aux messages planifiés, différés, ordonnés ou transactionnels, et sont transmis aux consommateurs aussi rapidement que possible.

Prérequis

Avant de commencer, assurez-vous d'avoir :

  • Installé le SDK Java. Pour plus d'informations, consultez la section Préparer l'environnement

  • Créé les ressources requises dans la console ApsaraMQ for RocketMQ : une instance, un topic (configuré avec le type de message Normal) et un groupe de consommateurs. Pour plus d'informations, consultez la section Créer des ressources

  • Récupéré votre paire AccessKey (AccessKey ID et AccessKey secret). Pour plus d'informations, consultez la section Créer une paire AccessKey

Fonctionnement

Les deux exemples suivent le même flux de travail :

  1. Créez un objet MQClient en indiquant l'endpoint HTTP et les identifiants AccessKey.

  2. Associez un producteur ou un consommateur à une instance et à un topic spécifiques.

  3. Envoyez ou interrogez les messages, puis fermez le client une fois l'opération terminée.

Chaque topic ne prend en charge qu'un seul type de message. Un topic créé pour les messages normaux ne peut ni envoyer ni recevoir de messages planifiés, différés, ordonnés ou transactionnels.

Envoyer un message

L'exemple suivant envoie quatre messages normaux vers un topic dans une boucle, en utilisant la publication synchrone. Si la méthode publishMessage se termine sans lever d'exception, le message a bien été envoyé.

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

Espace réservé Description Où le trouver
<your-http-endpoint> Endpoint HTTP de l'instance Page Instance Details > Section HTTP Endpoint
<your-topic> Nom du topic Console ApsaraMQ for RocketMQ
<your-instance-id> ID de l'instance. Définissez cette valeur sur null ou "" si l'instance n'utilise pas de namespace Page Instance Details
import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQProducer;
import com.aliyun.mq.http.model.TopicMessage;

import java.util.Date;

public class Producer {

    public static void main(String[] args) {
        MQClient mqClient = new MQClient(
                "<your-http-endpoint>",
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
        );

        final String topic = "<your-topic>";
        final String instanceId = "<your-instance-id>";

        // Get a producer for this topic and instance.
        MQProducer producer;
        if (instanceId != null && instanceId != "") {
            producer = mqClient.getProducer(instanceId, topic);
        } else {
            producer = mqClient.getProducer(topic);
        }

        try {
            for (int i = 0; i < 4; i++) {
                TopicMessage pubMsg = new TopicMessage(
                        "hello mq!".getBytes(),  // Message body
                        "A"                       // Message tag for filtering
                );
                // Custom property
                pubMsg.getProperties().put("a", String.valueOf(i));
                // Business key for message tracing (use an order ID, user ID, or similar identifier)
                pubMsg.setMessageKey("MessageKey");

                // Publish synchronously. No exception means success.
                TopicMessage pubResultMsg = producer.publishMessage(pubMsg);

                System.out.println(new Date() + " Send mq message success. Topic is:" + topic
                        + ", msgId is: " + pubResultMsg.getMessageId()
                        + ", bodyMD5 is: " + pubResultMsg.getMessageBodyMD5());
            }
        } catch (Throwable e) {
            System.out.println(new Date() + " Send mq message failed. Topic is:" + topic);
            e.printStackTrace();
        }

        mqClient.close();
    }
}

Points clés :

  • Identifiants issus des variables d'environnement : le code lit les valeurs ALIBABA_CLOUD_ACCESS_KEY_ID et ALIBABA_CLOUD_ACCESS_KEY_SECRET depuis les variables d'environnement. Définissez ces variables avant d'exécuter le code.

  • Tag de message : les tags permettent aux abonnés de filtrer les messages au sein d'un topic. Dans cet exemple, le tag est "A".

  • Clé de message : définissez la clé de message à l'aide d'un identifiant métier, tel qu'un ID de commande ou un ID utilisateur. Cela facilite la recherche et le suivi des messages dans la console ApsaraMQ for RocketMQ.

Recevoir et acquitter des messages

L'exemple suivant consomme les messages d'un topic dans une boucle utilisant le long polling. Après le traitement de chaque lot, il envoie un accusé de réception (ACK) au broker. Si le broker ne reçoit pas d'ACK avant l'expiration de l'intervalle de nouvelle tentative de livraison, il redistribue le message.

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

Espace réservé Description Où le trouver
<your-http-endpoint> Endpoint HTTP de l'instance Page Instance Details > Section HTTP Endpoint
<your-topic> Nom du topic Console ApsaraMQ for RocketMQ
<your-instance-id> ID de l'instance. Définissez cette valeur sur null ou "" si l'instance n'utilise pas de namespace Page Instance Details
<your-group-id> ID du groupe de consommateurs Console ApsaraMQ for RocketMQ
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(
                "<your-http-endpoint>",
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
        );

        final String topic = "<your-topic>";
        final String groupId = "<your-group-id>";
        final String instanceId = "<your-instance-id>";

        // Get a consumer for this topic, group, and instance.
        final MQConsumer consumer;
        if (instanceId != null && instanceId != "") {
            consumer = mqClient.getConsumer(instanceId, topic, groupId, null);
        } else {
            consumer = mqClient.getConsumer(topic, groupId);
        }

        // Poll for messages in a loop.
        // In production, use multiple threads for concurrent consumption.
        do {
            List<Message> messages = null;

            try {
                // Long polling: if no message is available, the request hangs
                // on the broker for up to the specified number of seconds.
                messages = consumer.consumeMessage(
                        3,  // Max messages per batch (up to 16)
                        3   // Long-polling timeout in seconds (up to 30)
                );
            } catch (Throwable e) {
                e.printStackTrace();
                try {
                    Thread.sleep(2000);
                } catch (InterruptedException e1) {
                    e1.printStackTrace();
                }
            }

            if (messages == null || messages.isEmpty()) {
                System.out.println(Thread.currentThread().getName() + ": no new message, continue!");
                continue;
            }

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

            // --- Acknowledge messages ---
            // Each message has a unique receipt handle that expires after
            // the delivery retry interval. ACK before it expires to prevent
            // redelivery.
            {
                List<String> handles = new ArrayList<String>();
                for (Message message : messages) {
                    handles.add(message.getReceiptHandle());
                }

                try {
                    consumer.ackMessage(handles);
                } catch (Throwable e) {
                    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);
    }
}

Points clés :

  • Long polling : le consommateur maintient la connexion ouverte sur le broker pendant la durée du délai d'attente spécifié (3 secondes dans cet exemple). Si un message arrive durant cette fenêtre, le broker répond immédiatement au lieu d'attendre le prochain cycle d'interrogation. Le délai maximal est de 30 secondes.

  • Taille du lot : chaque appel à consumeMessage renvoie jusqu'au nombre spécifié de messages (3 dans cet exemple, maximum 16).

  • Receipt handle : à chaque livraison, un message reçoit un nouveau receipt handle associé à un horodatage unique. Utilisez ce handle pour acquitter le message. Si le handle expire avant que l'ACK n'atteigne le broker, le message est redistribué.

  • Concurrence : cet exemple utilise un seul thread. Pour obtenir un débit plus élevé, consommez les messages via plusieurs threads en environnement de production.

Remarques d'utilisation

  • Un seul type par topic : un topic créé pour les messages normaux ne peut ni envoyer ni recevoir d'autres types de messages (planifiés, différés, ordonnés ou transactionnels). Créez des topics distincts pour chaque type de message.

  • Définir des clés de message pertinentes : utilisez des identifiants métiers (tels que des IDs de commande ou des IDs utilisateur) comme clé de message. Cela accélère la recherche et le suivi des messages dans la console.

  • Gérer correctement les échecs d'acquittement : si ackMessage lève une exception AckMessageException, enregistrez les handles ayant échoué ainsi que les codes d'erreur. Le broker redistribue automatiquement les messages non acquittés ; évitez donc de traiter les échecs d'ACK comme des erreurs fatales.

  • Libérer les ressources : appelez mqClient.close() lorsque le producteur ou le consommateur n'est plus nécessaire afin de libérer les connexions HTTP.

Étapes suivantes