Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive messages with the Java SDK

Dernière mise à jour :Aug 09, 2026

Une fois les ressources ApsaraMQ for RocketMQ créées (topic et groupe de consommateurs), intégrez la messagerie à votre application Java. Ce guide explique comment ajouter la dépendance du SDK Java 5.x, vous connecter à votre instance, envoyer des messages normaux et les consommer via un consommateur push ou un consommateur simple.

Prérequis

Vérifiez que vous disposez des éléments suivants :

Ajouter la dépendance du SDK

  1. Créez un projet Java dans votre IDE.

  2. Ajoutez la dépendance suivante au fichier pom.xml :

        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-client-java</artifactId>
            <version>5.0.7</version>
        </dependency>

Rassembler les paramètres de connexion

Avant d'écrire le code, récupérez les valeurs suivantes depuis la console ApsaraMQ for RocketMQ :

Paramètre Exemple Emplacement
endpoints rmq-cn-xxx.{regionId}.rmq.aliyuncs.com:8080 Page Instance Details > onglet Endpoints. Utilisez l'endpoint VPC pour un accès interne ou l'endpoint public pour un accès Internet. Consultez la rubrique Obtenir l'endpoint d'une instance.
topic normal_test Topic créé pour l'envoi et la réception de messages. Il doit exister préalablement. Consultez la rubrique Créer un topic.
consumerGroup GID_test Groupe de consommateurs créé. Il doit exister préalablement. Consultez la rubrique Créer un groupe de consommateurs.
InstanceId rmq-cn-xxx ID de votre instance ApsaraMQ for RocketMQ. Requis uniquement pour les instances serverless accessibles via Internet.
Nom d'utilisateur de l'instance 1XVg0hzgKm****** Page Access Control > onglet Intelligent Authentication. Consultez la rubrique Obtenir le nom d'utilisateur et le mot de passe d'une instance.
Mot de passe de l'instance ijSt8rEc45****** Même emplacement que le nom d'utilisateur.

Quand fournir les identifiants ?

La nécessité de fournir un nom d'utilisateur, un mot de passe et un namespace dépend de votre méthode d'accès :

Méthode d'accès Nom d'utilisateur et mot de passe Namespace (ID d'instance)
VPC (instance standard) Non requis. Le broker récupère automatiquement les identifiants depuis les informations VPC. Non requis.
VPC (instance serverless, authentification sans mot de passe activée) Non requis. Non requis.
VPC (instance serverless, authentification sans mot de passe désactivée) Requis. Non requis.
Internet (instance standard) Requis. Activez d'abord l'accès Internet sur l'instance. Non requis.
Internet (instance serverless) Requis. Requis. Définissez-le via .setNamespace("InstanceId").

Envoyer des messages

Créez un fichier ProducerExample.java dans votre projet et exécutez-le :

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 {
        // Instance endpoint (VPC endpoint for internal access, public endpoint for Internet access)
        String endpoints = "<your-endpoint>";
        // Topic (must be pre-created in the console)
        String topic = "<your-topic>";

        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
                .setEndpoints(endpoints)
                // Uncomment for serverless instances accessed over the Internet:
                // .setNamespace("<your-instance-id>")
                // Uncomment for Internet access or serverless VPC without authentication-free:
                // .setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"))
                .build();

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

        Message message = provider.newMessageBuilder()
                .setTopic(topic)
                .setKeys("messageKey")
                .setTag("messageTag")
                .setBody("messageBody".getBytes())
                .build();
        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println(sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

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

Espace réservé Description
<your-endpoint> Endpoint de l'instance, par exemple rmq-cn-xxx.cn-hangzhou.rmq.aliyuncs.com:8080
<your-topic> Nom du topic, par exemple normal_test
<your-instance-id> ID de l'instance (uniquement pour l'accès Internet serverless)
<instance-username> Nom d'utilisateur de l'instance (accès Internet ou VPC serverless sans authentification sans mot de passe)
<instance-password> Mot de passe de l'instance (accès Internet ou VPC serverless sans authentification sans mot de passe)
Remarque

L'appel à .setTopics() lors de l'initialisation du producteur valide les paramètres du topic en amont. Cette étape est facultative pour les messages normaux (validés dynamiquement à l'envoi), mais requise pour les messages transactionnels afin d'éviter les échecs de l'API de requête.

Recevoir des messages

ApsaraMQ for RocketMQ propose deux types de consommateurs :

Consommateur push Consommateur simple
Fonctionnement Le SDK livre les messages à un callback (écouteur de messages). Votre application extrait les messages et accuse explicitement réception de chacun d'eux.
Concurrence Gérée par le SDK. Gérée par votre application.
Flexibilité Plus faible. Le SDK encapsule le flux de consommation. Plus élevée. Les opérations atomiques permettent de créer des flux de travail personnalisés.
Cas d'utilisation idéal La plupart des cas nécessitant le traitement des messages entrants. Scénarios exigeant un contrôle fin de l'extraction, du traitement ou de l'accusé de réception.

Privilégiez un consommateur push, sauf si vous avez besoin d'un contrôle personnalisé du flux de consommation.

Consommateur push

Créez un fichier PushConsumerExample.java et exécutez-le :

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);

    private PushConsumerExample() {
    }

    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();
        ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
                .setEndpoints(endpoints)
                // Uncomment for serverless instances accessed over the Internet:
                // .setNamespace("<your-instance-id>")
                // Uncomment for Internet access or serverless VPC without authentication-free:
                // .setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"))
                .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("Consume Message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();

        // Keep the consumer running
        Thread.sleep(Long.MAX_VALUE);

        // To shut down the consumer gracefully, call:
        // pushConsumer.close();
    }
}

Consommateur simple

Créez un fichier SimpleConsumerExample.java et exécutez-le :

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);

    private SimpleConsumerExample() {
    }

    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();
        ClientConfiguration clientConfiguration = ClientConfiguration.newBuilder()
                .setEndpoints(endpoints)
                // Uncomment for serverless instances accessed over the Internet:
                // .setNamespace("<your-instance-id>")
                // Uncomment for Internet access or serverless VPC without authentication-free:
                // .setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"))
                .build();

        // Long polling timeout
        Duration awaitDuration = Duration.ofSeconds(10);
        // Subscribe to all messages in the topic (tag = "*")
        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();

        // Maximum messages per pull
        int maxMessageNum = 16;
        // Messages stay invisible to other consumers for this duration after being received
        Duration invisibleDuration = Duration.ofSeconds(10);

        // Poll for messages in a loop. For real-time consumption, use multiple threads.
        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 {
                    // ACK each message to commit the consumption result to the broker
                    consumer.ack(message);
                    System.out.println("Message is acknowledged successfully, messageId= " + messageId);
                } catch (Throwable t) {
                    t.printStackTrace();
                }
            }
        }
        // To shut down the consumer gracefully, call:
        // consumer.close();
    }
}

Vérifier la livraison des messages

Consultez l'état de la livraison dans la console ApsaraMQ for RocketMQ :

  1. Connectez-vous à la console ApsaraMQ for RocketMQ.

  2. Sur la page Instances, cliquez sur le nom de votre instance.

  3. Dans le volet de navigation de gauche, cliquez sur Message Tracing.

Versions du SDK requises pour l'accès Internet serverless

L'accès à une instance ApsaraMQ for RocketMQ serverless via Internet nécessite une version minimale du SDK. Remplacez InstanceId dans les exemples par votre ID d'instance réel.

SDK for Java 5.x (rocketmq-client-java)

Version minimale : 5.0.6

Définissez le namespace dans ClientConfiguration :

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

SDK for Java 5.x (rocketmq-client)

Version minimale : 5.2.0

Définissez le namespace séparément sur le producteur et le consommateur :

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

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

TCP client SDK for Java 1.x

Version minimale : 1.9.0.Final

Définissez le namespace via les propriétés :

properties.setProperty(PropertyKeyConst.Namespace, "InstanceId");

Étapes suivantes