Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Normal messages

Dernière mise à jour :Aug 09, 2026

Les messages normaux sont le type de message par défaut dans ApsaraMQ for RocketMQ. Contrairement aux messages ordonnés, planifiés, différés et transactionnels, ils ne disposent d'aucune fonctionnalité spéciale. Ils permettent une communication asynchrone et découplée entre les producteurs et les consommateurs. Utilisez des messages normaux lorsque votre application exige une livraison fiable sans imposer un ordre strict ni une distribution programmée.

Quand utiliser les messages normaux

Les messages normaux conviennent à tout scénario qui privilégie la fiabilité de la transmission à l'ordre de traitement. Voici quelques exemples courants : découplage des microservices, intégration de données et architectures orientées événements.

  • Découplage des microservices : Un service amont encapsule la prise de commande et le paiement dans un message normal indépendant, puis l'envoie au broker ApsaraMQ for RocketMQ. Les services aval s'abonnent indépendamment et traitent le message selon leur propre logique. Chaque message est autonome : il ne dépend d'aucun autre message.

    Asynchronous decoupling between an order system and downstream services

  • Intégration de données : Un composant d'instrumentation collecte les journaux d'opération des applications frontales et les transfère vers ApsaraMQ for RocketMQ. Chaque message correspond à un enregistrement de journal discret. Le broker stocke et distribue ces enregistrements aux systèmes de stockage aval sans les transformer. Les applications backend traitent ensuite les tâches associées.

    Log data pipeline from frontend applications to downstream storage

Cycle de vie des messages

Un message normal traverse cinq étapes entre sa création et sa suppression :

Lifecycle of a normal message

Étape Description
Initialisé Le producteur construit le message et le prépare pour l'envoi.
Prêt Le broker accepte le message. Celui-ci devient alors visible et disponible pour les consommateurs.
En cours Un consommateur récupère le message et commence son traitement. Si le broker ne reçoit aucun accusé de réception dans le délai imparti, il réessaie la livraison.
Accusé Le consommateur valide le résultat de la consommation. Le broker marque le message comme consommé, mais ne le supprime pas immédiatement.
Supprimé Le message est supprimé à l'expiration de la période de rétention ou lorsque l'espace de stockage vient à manquer. La suppression suit une stratégie rotative qui retire d'abord les messages les plus anciens du fichier physique. Pour plus de détails, consultez la rubrique Stockage et nettoyage des messages.

Par défaut, ApsaraMQ for RocketMQ conserve tous les messages, même après accusé de réception. Les consommateurs peuvent reconsommer tout message qui n'a pas encore été supprimé.

Exemples de code (Java)

Les exemples suivants utilisent le SDK Apache RocketMQ 5.x pour envoyer et consommer des messages normaux. Chaque exemple configure une connexion client et exécute une seule opération. Remplacez les espaces réservés par vos valeurs réelles avant d'exécuter le code.

Pour obtenir la référence complète du SDK et des exemples dans d'autres langages, consultez la section SDK Apache RocketMQ 5.x.

Avant de commencer

Préparez les ressources suivantes dans la console ApsaraMQ for RocketMQ :

Ressource Où la trouver
Endpoint Page Instance Details > onglet Endpoints. Utilisez l'endpoint VPC pour accéder depuis des instances Elastic Compute Service (ECS) situées dans le même Virtual Private Cloud (VPC). Utilisez l'endpoint public pour un accès Internet (nécessite l'activation de la fonctionnalité d'accès Internet).
Topic Créez un topic avec MessageType défini sur Normal. Les messages normaux ne peuvent être envoyés qu'à des topics de ce type.
Consumer group Requis pour les exemples de consommation. Créez-en un dans la console.
Credentials (endpoint public uniquement) Nom d'utilisateur et mot de passe disponibles sur la page Access Control > onglet Intelligent Authentication. Non requis pour l'accès VPC, car le broker s'authentifie automatiquement. Pour les instances serverless accessibles via un VPC avec la fonctionnalité d'authentification gratuite dans les VPC activée, les identifiants ne sont pas non plus requis.

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

Espace réservé Description Exemple
<your-endpoint> Endpoint de l'instance rmq-cn-xxx.cn-hangzhou.rmq.aliyuncs.com:8080
<your-topic> Nom du topic normal-topic-01
<your-consumer-group> Nom du consumer group GID_normal_consumer
<your-username> Nom d'utilisateur de l'instance (endpoint public uniquement) --
<your-password> Mot de passe de l'instance (endpoint public uniquement) --

Exemple de code

Envoyer un message normal

package doc;

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.apis.producer.Producer;
import org.apache.rocketmq.client.apis.producer.SendReceipt;

public class ProducerExample {
    public static void main(String[] args) throws ClientException {
        // Specify the instance endpoint.
        String endpoints = "<your-endpoint>";
        String topic = "<your-topic>";

        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // Uncomment the following line for public endpoint access:
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        ClientConfiguration configuration = builder.build();

        Producer producer = provider.newProducerBuilder()
                .setTopics(topic)
                .setClientConfiguration(configuration)
                .build();

        // Build a normal message with a key and tag for filtering and tracing.
        Message message = provider.newMessageBuilder()
                .setTopic(topic)
                .setKeys("messageKey")
                .setTag("messageTag")
                .setBody("messageBody".getBytes())
                .build();

        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println("Send success. Message ID: " + sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

Sortie attendue :

Send success. Message ID: 01BE0A3600F5762D0455025C2D1A0000

Consommer avec un push consumer

Un push consumer reçoit les messages via un callback MessageListener. Le broker pousse automatiquement les messages vers le consommateur.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.ConsumeResult;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.PushConsumer;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.util.Collections;

public class PushConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(PushConsumerExample.class);

    public static void main(String[] args) throws ClientException, IOException, InterruptedException {
        String endpoints = "<your-endpoint>";
        String topic = "<your-topic>";
        String consumerGroup = "<your-consumer-group>";

        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // Uncomment the following line for public endpoint access:
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        ClientConfiguration clientConfiguration = builder.build();

        // Subscribe to all messages in the topic (tag = "*").
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);

        PushConsumer pushConsumer = provider.newPushConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .setMessageListener(messageView -> {
                    System.out.println("Received message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();

        // Keep the process running to continue receiving messages.
        Thread.sleep(Long.MAX_VALUE);

        // Close the consumer when it is no longer needed.
        // pushConsumer.close();
    }
}

Sortie attendue :

Received message: MessageView{messageId=01BE0A3600F5762D0455025C2D1A0000, topic=<your-topic>, ...}

Consommer avec un simple consumer

Un simple consumer extrait explicitement les messages via receive() et accuse réception de chaque message avec ack(). Ce mode vous offre un contrôle total sur le rythme de consommation.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.SimpleConsumer;
import org.apache.rocketmq.client.apis.message.MessageId;
import org.apache.rocketmq.client.apis.message.MessageView;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.time.Duration;
import java.util.Collections;
import java.util.List;

public class SimpleConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(SimpleConsumerExample.class);

    public static void main(String[] args) throws ClientException, IOException {
        String endpoints = "<your-endpoint>";
        String topic = "<your-topic>";
        String consumerGroup = "<your-consumer-group>";

        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // Uncomment the following line for public endpoint access:
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<your-username>", "<your-password>"));
        ClientConfiguration clientConfiguration = builder.build();

        Duration awaitDuration = Duration.ofSeconds(10);
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);

        SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setAwaitDuration(awaitDuration)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .build();

        int maxMessageNum = 16;
        Duration invisibleDuration = Duration.ofSeconds(10);

        // Pull and acknowledge messages in a loop.
        // For real-time consumption, use multiple threads to pull concurrently.
        while (true) {
            final List<MessageView> messages = consumer.receive(maxMessageNum, invisibleDuration);
            messages.forEach(messageView -> {
                System.out.println("Received message: " + messageView);
            });
            for (MessageView message : messages) {
                final MessageId messageId = message.getMessageId();
                try {
                    consumer.ack(message);
                    System.out.println("Acked message: " + messageId);
                } catch (Throwable t) {
                    t.printStackTrace();
                }
            }
        }

        // Close the consumer when it is no longer needed.
        // consumer.close();
    }
}

Sortie attendue :

Received message: MessageView{messageId=01BE0A3600F5762D0455025C2D1A0000, topic=<your-topic>, ...}
Acked message: 01BE0A3600F5762D0455025C2D1A0000

Bonnes pratiques

Définir une clé de message globalement unique pour chaque message

Attribuez un identifiant métier spécifique, tel qu'un ID de commande ou un ID utilisateur, en tant que clé de message. Cela vous permet de rechercher des messages individuels et leurs traces dans la console ApsaraMQ for RocketMQ.

Message message = provider.newMessageBuilder()
        .setTopic(topic)
        .setKeys("order-12345")   // Use a unique business identifier.
        .setTag("payment")
        .setBody(payload)
        .build();

Choisir le bon mode de consommation

Mode Fonctionnement Recommandé pour
Push consumer Le broker pousse les messages vers un callback MessageListener. Le traitement événementiel à faible latence, où le consommateur traite chaque message dès sa réception.
Simple consumer L'application appelle receive() pour extraire un lot, puis ack() chaque message après traitement. Les charges de travail nécessitant un contrôle manuel du flux ou un traitement par lots.